Pour exécuter des requêtes Spark SQL de manière interactive, spécifiez un groupe de ressources Spark Interactive. Ce groupe de ressources s'adapte automatiquement dans une plage définie afin de répondre à vos besoins d'analyse interactive tout en réduisant les coûts. Cette rubrique explique comment effectuer une analyse interactive avec Spark SQL à l'aide de la console, de Hive JDBC, de PyHive, de Beeline, de DBeaver et d'autres outils clients.
Prérequis
Un cluster AnalyticDB for MySQL Enterprise Edition, Basic Edition ou Data Lakehouse Edition est créé.
Un bucket Object Storage Service (OSS) est créé dans la même région que le cluster AnalyticDB for MySQL.
-
Un compte de base de données est créé pour le cluster AnalyticDB for MySQL.
Si vous utilisez un compte Alibaba Cloud, il vous suffit de créer un compte privilégié.
Si vous utilisez un utilisateur Resource Access Management (RAM), vous devez créer un compte privilégié et un compte standard puis associer le compte standard à l'utilisateur RAM.
Un environnement de développement Java 8 et Python 3.9 doit être installé pour exécuter des clients tels que des applications Java, des applications Python et Beeline.
L'adresse IP de votre client doit figurer dans la liste d'autorisation du cluster AnalyticDB for MySQL.
Remarques sur l'utilisation
Si un groupe de ressources Spark Interactive est arrêté, le cluster le redémarre lorsque vous exécutez la première requête Spark SQL. La première requête peut être mise en file d'attente pendant le démarrage.
Spark ne peut ni lire ni écrire dans les bases de données INFORMATION_SCHEMA et MYSQL. N'utilisez pas ces bases de données comme base de données de connexion initiale.
Assurez-vous que le compte de base de données utilisé pour soumettre les tâches Spark SQL peut accéder à la base de données cible. Sinon, la requête échouera.
Préparation
-
Obtenez l'endpoint du groupe de ressources Spark Interactive.
Connectez-vous à la console AnalyticDB for MySQL. Dans le coin supérieur gauche de la console, sélectionnez une région. Dans le volet de navigation de gauche, cliquez sur Clusters. Recherchez le cluster que vous souhaitez gérer et cliquez sur son ID.
Dans le volet de navigation de gauche, choisissez , puis cliquez sur l'onglet Resource Groups.
-
Recherchez le groupe de ressources et cliquez sur Details dans la colonne Actions pour afficher l'endpoint interne et l'endpoint public. Vous pouvez cliquer sur l'icône
située à côté d'un endpoint pour le copier, ou cliquer sur l'icône
située entre les parenthèses du numéro de Port pour copier la chaîne de connexion JDBC.Dans les cas suivants, vous devez cliquer sur Apply for Endpoint à côté de Public Address pour demander manuellement un endpoint public.
L'outil client utilisé pour soumettre les tâches Spark SQL est déployé sur votre machine locale ou sur un serveur externe.
L'outil client utilisé pour soumettre les tâches Spark SQL est déployé sur une instance ECS, et l'instance ECS et le cluster AnalyticDB for MySQL ne se trouvent pas dans le même VPC.
Les informations de connexion incluent également des champs tels que le port public et le port VPC (par défaut :
10000), l'ID VPC, l'ID vSwitch, la classe de pilote (org.apache.hive.jdbc.HiveDriver) et l'URL de téléchargement du pilote.
Analyse interactive
Console
Si vous utilisez un HiveMetastore géré par vos soins, créez une base de données nommée default dans AnalyticDB for MySQL et sélectionnez-la comme base de données lors de l'exécution des tâches Spark SQL dans la console.
Connectez-vous à la console AnalyticDB for MySQL. Dans le coin supérieur gauche de la console, sélectionnez une région. Dans le volet de navigation de gauche, cliquez sur Clusters. Recherchez le cluster que vous souhaitez gérer et cliquez sur son ID.
Dans le volet de navigation de gauche, choisissez .
-
Sélectionnez le moteur Spark et le groupe de ressources Spark Interactive créé, puis exécutez l'instruction Spark SQL suivante :
SHOW DATABASES;
SDK
Lorsque vous exécutez des instructions Spark SQL à l'aide d'un SDK, les résultats de la requête sont écrits sous forme de fichiers dans un bucket OSS spécifié. Vous pouvez ensuite interroger les données dans la console OSS ou télécharger les fichiers de résultats sur votre ordinateur. L'exemple suivant montre comment appeler le SDK en Python.
-
Exécutez la commande suivante pour installer le SDK.
pip install alibabacloud-adb20211201 -
Exécutez les commandes suivantes pour installer les dépendances.
pip install oss2 pip install loguru -
Connectez-vous au cluster et exécutez des instructions Spark SQL.
# coding: utf-8 import csv import json import time from io import StringIO import oss2 from alibabacloud_adb20211201.client import Client from alibabacloud_adb20211201.models import ExecuteSparkWarehouseBatchSQLRequest, ExecuteSparkWarehouseBatchSQLResponse, \ GetSparkWarehouseBatchSQLRequest, GetSparkWarehouseBatchSQLResponse, \ ListSparkWarehouseBatchSQLRequest, CancelSparkWarehouseBatchSQLRequest, ListSparkWarehouseBatchSQLResponse from alibabacloud_tea_openapi.models import Config from loguru import logger def build_sql_config(oss_location, spark_sql_runtime_config: dict = None, file_format = "CSV", output_partitions = 1, sep = "|"): """ Builds the configuration for an AnalyticDB for MySQL SQL execution. :param oss_location: The OSS path to store the SQL execution results. :param spark_sql_runtime_config: The native Spark SQL configuration properties. :param file_format: The file format of the SQL execution results. Default value: CSV. :param output_partitions: The number of partitions for the SQL execution results. If you need to output a large result set, you must increase this value to avoid creating a single oversized file. :param sep: The separator for CSV files. This parameter is ignored for non-CSV files. :return: The configuration for the SQL execution. """ if oss_location is None: raise ValueError("oss_location is required") if not oss_location.startswith("oss://"): raise ValueError("oss_location must start with oss://") if file_format != "CSV" and file_format != "PARQUET" and file_format != "ORC" and file_format != "JSON": raise ValueError("file_format must be CSV, PARQUET, ORC or JSON") runtime_config = { # sql output config "spark.adb.sqlOutputFormat": file_format, "spark.adb.sqlOutputPartitions": output_partitions, "spark.adb.sqlOutputLocation": oss_location, # csv config "sep": sep } if spark_sql_runtime_config: runtime_config.update(spark_sql_runtime_config) return runtime_config def execute_sql(client: Client, dbcluster_id: str, resource_group_name: str, query: str, limit = 10000, runtime_config: dict = None, schema="default" ): """ Runs an SQL statement in a Spark Interactive resource group. :param client: The Alibaba Cloud client. :param dbcluster_id: The ID of the cluster. :param resource_group_name: The resource group of the cluster. This must be a Spark Interactive resource group. :param schema: The name of the default database for the SQL execution. If you do not specify this parameter, the default value is used. :param limit: The number of rows to return for the SQL execution result. :param query: The SQL statement to run. Use semicolons (;) to separate multiple SQL statements. :return: """ # Assemble the request body. req = ExecuteSparkWarehouseBatchSQLRequest() # The cluster ID. req.dbcluster_id = dbcluster_id # The name of the resource group. req.resource_group_name = resource_group_name # The timeout period for the SQL execution. req.execute_time_limit_in_seconds = 3600 # The name of the database where the SQL statement is run. req.schema = schema # The SQL query or statements. req.query = query # The number of result rows to return. req.execute_result_limit = limit if runtime_config: # The configuration for the SQL execution. req.runtime_config = json.dumps(runtime_config) # Submits the SQL statement and returns the query ID. resp: ExecuteSparkWarehouseBatchSQLResponse = client.execute_spark_warehouse_batch_sql(req) logger.info("Query execute submitted: {}", resp.body.data.query_id) return resp.body.data.query_id def get_query_state(client, query_id): """ Queries the execution status of an SQL statement. :param client: The Alibaba Cloud client. :param query_id: The ID of the SQL execution. :return: The execution status and result of the SQL statement. """ req = GetSparkWarehouseBatchSQLRequest(query_id=query_id) resp: GetSparkWarehouseBatchSQLResponse = client.get_spark_warehouse_batch_sql(req) logger.info("Query state: {}", resp.body.data.query_state) return resp.body.data.query_state, resp def list_history_query(client, db_cluster, resource_group_name, page_num): """ Queries the history of SQL statements run in a Spark Interactive resource group. :param client: The Alibaba Cloud client. :param db_cluster: The ID of the cluster. :param resource_group_name: The name of the resource group. :param page_num: The page number for paginated queries. :return: Specifies whether more pages are available. If queries exist, you can proceed to the next page. """ req = ListSparkWarehouseBatchSQLRequest(dbcluster_id=db_cluster, resource_group_name=resource_group_name, page_number = page_num) resp: ListSparkWarehouseBatchSQLResponse = client.list_spark_warehouse_batch_sql(req) # If no SQL statement is found, return True. Otherwise, return True. The default is 10 entries per page. if resp.body.data.queries is None: return True # Print the queried SQL statements. for query in resp.body.data.queries: logger.info("Query ID: {}, State: {}", query.query_id, query.query_state) logger.info("Total queries: {}", len(resp.body.data.queries)) return len(resp.body.data.queries) < 10 def list_csv_files(oss_client, dir): for obj in oss_client.list_objects_v2(dir).object_list: if obj.key.endswith(".csv"): logger.info(f"reading {obj.key}") # read oss file content csv_content = oss_client.get_object(obj.key).read().decode('utf-8') csv_reader = csv.DictReader(StringIO(csv_content)) # Print the CSV content for row in csv_reader: print(row) if __name__ == '__main__': logger.info("ADB Spark Batch SQL Demo") # Replace with your AccessKey ID. _ak = "LTAI****************" # Replace with your AccessKey secret. _sk = "yourAccessKeySecret" # Replace with the actual region ID. _region= "cn-shanghai" # Replace with your cluster ID. _db = "amv-uf6485635f****" # Replace with the name of your resource group. _rg_name = "testjob" # client config client_config = Config( # Your Alibaba Cloud AccessKey ID. access_key_id=_ak, # Your Alibaba Cloud AccessKey secret. access_key_secret=_sk, # The endpoint of the AnalyticDB for MySQL service. # adb.ap-southeast-1.aliyuncs.com is the endpoint of the service in the China (Singapore) region. # adb-vpc.ap-southeast-1.aliyuncs.com is used in VPC scenarios. endpoint=f"adb.{_region}.aliyuncs.com" ) # Create an Alibaba Cloud client. _client = Client(client_config) # The configuration for the SQL execution. _spark_sql_runtime_config = { "spark.sql.shuffle.partitions": 1000, "spark.sql.autoBroadcastJoinThreshold": 104857600, "spark.sql.sources.partitionOverwriteMode": "dynamic", "spark.sql.sources.partitionOverwriteMode.dynamic": "dynamic" } _config = build_sql_config(oss_location="oss://testBucketName/sql_result", spark_sql_runtime_config = _spark_sql_runtime_config) # The SQL statement to run. _query = """ SHOW DATABASES; SELECT 100; """ _query_id = execute_sql(client = _client, dbcluster_id=_db, resource_group_name=_rg_name, query=_query, runtime_config=_config) logger.info(f"Run query_id: {_query_id} for SQL {_query}.\n Waiting for result...") # Wait for the SQL execution to complete. current_ts = time.time() while True: query_state, resp = get_query_state(_client, _query_id) """ The query_state can be one of the following: - PENDING: The query is queued, which can occur while the Spark Interactive resource group is starting. - SUBMITTED: The query is submitted to the Spark Interactive resource group. - RUNNING: The SQL statement is executing. - FINISHED: The SQL execution completed successfully. - FAILED: The SQL execution failed. - CANCELED: The SQL execution is canceled. """ if query_state == "FINISHED": logger.info("query finished success") break elif query_state == "FAILED": # Print the failure information. logger.error("Error Info: {}", resp.body.data) exit(1) elif query_state == "CANCELED": # Print the cancellation information. logger.error("query canceled") exit(1) else: time.sleep(2) if time.time() - current_ts > 600: logger.error("query timeout") # If the execution time exceeds 10 minutes, cancel the SQL execution. _client.cancel_spark_warehouse_batch_sql(CancelSparkWarehouseBatchSQLRequest(query_id=_query_id)) exit(1) # A query can contain multiple statements. The following loop processes each statement. for stmt in resp.body.data.statements: logger.info( f"statement_id: {stmt.statement_id}, result location: {stmt.result_uri}") # Sample code to view the results. _bucket = stmt.result_uri.split("oss://")[1].split("/")[0] _dir = stmt.result_uri.replace(f"oss://{_bucket}/", "").replace("//", "/") oss_client = oss2.Bucket(oss2.Auth(client_config.access_key_id, client_config.access_key_secret), f"oss-{_region}.aliyuncs.com", _bucket) list_csv_files(oss_client, _dir) # Query all SQL statements run in the Spark Interactive resource group. You can perform paginated queries. logger.info("List all history query") page_num = 1 no_more_page = list_history_query(_client, _db, _rg_name, page_num) while no_more_page: logger.info(f"List page {page_num}") page_num += 1 no_more_page = list_history_query(_client, _db, _rg_name, page_num)Paramètres :
_ak : l'AccessKey ID de votre compte Alibaba Cloud ou d'un utilisateur RAM disposant des autorisations d'accès sur AnalyticDB for MySQL. Pour savoir comment obtenir un AccessKey ID et un AccessKey Secret, consultez la section Comptes et autorisations.
_sk : l'AccessKey Secret de votre compte Alibaba Cloud ou d'un utilisateur RAM disposant des autorisations d'accès sur AnalyticDB for MySQL. Pour savoir comment obtenir un AccessKey ID et un AccessKey Secret, consultez la section Comptes et autorisations.
_region : l'ID de la région où réside votre cluster AnalyticDB for MySQL.
_db : l'ID du cluster AnalyticDB for MySQL.
_rg_name : le nom du groupe de ressources Spark Interactive.
-
oss_location (facultatif) : le chemin OSS où sont stockés les fichiers de résultats de requête.
Si vous ne spécifiez pas ce paramètre, vous ne pourrez consulter que les cinq premières lignes du résultat de la requête dans le journal Log de la page .
Applications
Hive JDBC
-
Configurez la dépendance Maven dans le fichier pom.xml.
<dependency> <groupId>org.apache.hive</groupId> <artifactId>hive-jdbc</artifactId> <version>2.3.9</version> </dependency> -
Établissez une connexion et exécutez des instructions Spark SQL.
public class java { public static void main(String[] args) throws Exception { Class.forName("org.apache.hive.jdbc.HiveDriver"); String url = "<JDBC-connection-string>"; Connection con = DriverManager.getConnection(url, "<username>", "<password>"); Statement stmt = con.createStatement(); ResultSet tables = stmt.executeQuery("show tables"); List<String> tbls = new ArrayList<>(); while (tables.next()) { System.out.println(tables.getString("tableName")); tbls.add(tables.getString("tableName")); } } }Paramètres :
Chaîne de connexion JDBC : la chaîne de connexion JDBC du groupe de ressources Spark Interactive obtenue dans la section Préparation. Remplacez default par le nom de la base de données à laquelle vous souhaitez vous connecter.
Nom d'utilisateur : le compte AnalyticDB for MySQL pour AnalyticDB for MySQL.
Mot de passe : le mot de passe du compte AnalyticDB for MySQL pour AnalyticDB for MySQL.
PyHive
-
Installez le client PyHive.
pip install pyhive -
Établissez une connexion et exécutez des instructions Spark SQL.
from pyhive import hive from TCLIService.ttypes import TOperationState cursor = hive.connect( host='<endpoint>', port=<port>, username='<resource_group_name>/<username>', password='<password>', auth='CUSTOM' ).cursor() cursor.execute('show tables') status = cursor.poll().operationState while status in (TOperationState.INITIALIZED_STATE, TOperationState.RUNNING_STATE): logs = cursor.fetch_logs() for message in logs: print(message) # If needed, an asynchronous query can be cancelled at any time with: # cursor.cancel() status = cursor.poll().operationState print(cursor.fetchall())Paramètres :
Endpoint : l'endpoint du groupe de ressources Spark Interactive que vous avez obtenu dans la section Preparation.
Port : le port du groupe de ressources Spark Interactive, qui est 10000.
Nom du groupe de ressources : le nom du groupe de ressources Spark Interactive.
Nom d'utilisateur : le AnalyticDB for MySQL pour AnalyticDB for MySQL.
Mot de passe : le mot de passe du AnalyticDB for MySQL pour AnalyticDB for MySQL.
Clients
Outre les clients Beeline, DBeaver, DBVisualizer et DataGrip décrits dans cette rubrique, vous pouvez également effectuer des analyses interactives à l'aide d'outils d'ordonnancement de workflows tels que Airflow, Azkaban et DolphinScheduler.
Beeline
-
Connectez-vous au groupe de ressources Spark Interactive.
Utilisez le format de commande suivant :
!connect <JDBC-connection-string> <username> <password>Chaîne de connexion JDBC : la chaîne de connexion JDBC du groupe de ressources Spark Interactive que vous avez obtenue dans la section Preparation. Remplacez default par le nom de la base de données à laquelle vous souhaitez vous connecter.
Nom d'utilisateur : le AnalyticDB for MySQL pour AnalyticDB for MySQL.
Mot de passe : le mot de passe du AnalyticDB for MySQL pour AnalyticDB for MySQL.
Exemple :
!connect jdbc:hive2://amv-bp1c3em7b2e****-spark.ads.aliyuncs.com:10000/adb_test spark_resourcegroup/AdbSpark14**** Spark23****Une connexion réussie renvoie la sortie suivante :
Connected to: Spark SQL (version 3.2.0) Driver: Hive JDBC (version 2.3.9) Transaction isolation: TRANSACTION_REPEATABLE_READ -
Exécutez une instruction Spark SQL.
SHOW TABLES;
DBeaver
Ouvrez le client DBeaver et choisissez .
Sur la page Connect to a database, sélectionnez Apache Spark et cliquez sur Next.
-
Configurez les paramètres Hadoop/Apache Spark connection settings comme suit :
Paramètre
Description
Méthode de connexion
Sélectionnez URL.
URL JDBC
Saisissez la chaîne de connexion JDBC que vous avez obtenue dans la section Preparation.
ImportantRemplacez
defaultdans la chaîne de connexion par le nom de votre base de données.Nom d'utilisateur
Le AnalyticDB for MySQL pour AnalyticDB for MySQL.
Mot de passe
Le mot de passe du AnalyticDB for MySQL pour AnalyticDB for MySQL.
-
Après avoir configuré les paramètres, cliquez sur Test Connection.
ImportantLors du premier test de connexion, DBeaver vous invite à télécharger les pilotes requis. Cliquez sur Download pour les télécharger.
Une fois le test de connexion réussi, cliquez sur Finish.
Dans l'onglet Database Navigator, développez la source de données et cliquez sur la base de données.
-
Dans l'éditeur de code à droite, saisissez une instruction SQL et cliquez sur l'icône
pour l'exécuter.SHOW TABLES;+-----------+-----------+-------------+ | namespace | tableName | isTemporary | +-----------+-----------+-------------+ | db | test | [] | +-----------+-----------+-------------+
DBVisualizer
Ouvrez le client DBVisualizer et choisissez .
Sur la page Driver Manager, sélectionnez Hive et cliquez sur l'icône
.-
Dans l'onglet Driver Settings, configurez les paramètres suivants :
Paramètre
Description
Name
Un nom personnalisé pour la source de données Hive.
Format d'URL
Saisissez la chaîne de connexion JDBC que vous avez obtenue dans la section Preparation.
ImportantRemplacez
defaultdans la chaîne de connexion par le nom de votre base de données.Classe de pilote
Sélectionnez org.apache.hive.jdbc.HiveDriver.
ImportantAprès avoir configuré les paramètres, cliquez sur Start Download pour télécharger le pilote.
Une fois le pilote téléchargé, choisissez .
-
Dans la boîte de dialogue Create Database Connection from Database URL, configurez les paramètres décrits dans le tableau suivant.
Paramètre
Description
URL de la base de données
Saisissez la chaîne de connexion JDBC que vous avez obtenue dans la section Preparation.
ImportantRemplacez
defaultdans la chaîne de connexion par le nom de votre base de données.Classe de pilote
Sélectionnez la source de données Hive que vous avez créée à l'étape 3.
-
Sur la page Connection, configurez les paramètres de connexion suivants et cliquez sur Connect.
Paramètre
Description
Name
Par défaut, ce paramètre est défini sur le nom de la source de données Hive que vous avez créée à l'étape 3. Vous pouvez personnaliser le nom.
Notes
Saisissez des remarques.
Type de pilote
Sélectionnez Hive.
URL de la base de données
Saisissez la chaîne de connexion JDBC que vous avez obtenue dans la section Preparation.
ImportantRemplacez
defaultdans la chaîne de connexion par le nom de votre base de données.Identifiant utilisateur de la base de données
Le AnalyticDB for MySQL pour AnalyticDB for MySQL.
Mot de passe de la base de données
Le mot de passe du AnalyticDB for MySQL pour AnalyticDB for MySQL.
RemarqueLaissez les autres paramètres avec leurs valeurs par défaut.
Une fois la connexion établie, dans l'onglet Database, développez la source de données et cliquez sur la base de données.
-
Dans l'éditeur de code à droite, saisissez une instruction SQL et cliquez sur l'icône
pour l'exécuter.SHOW TABLES;+-----------+-----------+-------------+ | namespace | tableName | isTemporary | +-----------+-----------+-------------+ | db | test | false | +-----------+-----------+-------------+
DataGrip
Ouvrez le client DataGrip, choisissez et créez un projet.
-
Ajoutez une source de données.
Cliquez sur l'icône
et choisissez .-
Dans la boîte de dialogue Data Sources and Drivers qui s'affiche, configurez les paramètres suivants et cliquez sur OK.
Définissez Driver sur Apache Spark et Authentication sur User & Password.
Paramètre
Description
Name
Le nom de la source de données, que vous pouvez personnaliser. Cette rubrique utilise
adbtestà titre d'exemple.Host
Saisissez la chaîne de connexion JDBC que vous avez obtenue dans la section Preparation.
ImportantRemplacez
defaultdans la chaîne de connexion par le nom de votre base de données.Port
Le port du groupe de ressources Spark Interactive, qui est 10000.
User
Le AnalyticDB for MySQL pour AnalyticDB for MySQL.
Password
Le mot de passe du AnalyticDB for MySQL pour AnalyticDB for MySQL.
Schema
Le nom de la base de données dans le cluster AnalyticDB for MySQL.
-
Exécutez des instructions Spark SQL.
Dans la liste des sources de données, cliquez avec le bouton droit sur la source de données que vous avez créée à l'étape 2 et choisissez .
-
Dans le panneau Console qui s'affiche, exécutez une instruction Spark SQL.
SHOW TABLES;
Outils BI
Vous pouvez effectuer des analyses interactives à l'aide d'outils BI tels que Redash, Power BI et Metabase.