Para executar consultas Spark SQL interativamente, especifique um grupo de recursos Spark Interactive. O grupo de recursos dimensiona automaticamente dentro de um intervalo definido para atender às necessidades de análise interativa e reduzir custos. Este tópico descreve como realizar análises interativas com Spark SQL usando o console, Hive JDBC, PyHive, Beeline, DBeaver e outras ferramentas cliente.
Pré-requisitos
Um cluster do AnalyticDB for MySQL Enterprise Edition, Basic Edition ou Data Lakehouse Edition criado.
Um bucket do Object Storage Service (OSS) criado na mesma região do cluster do AnalyticDB for MySQL.
-
Uma conta de banco de dados criada para o cluster do AnalyticDB for MySQL.
Se você usar uma conta Alibaba Cloud, basta criar uma conta privilegiada.
Se você usar um usuário do Resource Access Management (RAM), crie uma conta privilegiada e uma conta padrão e associe a conta padrão ao usuário RAM.
Ambientes de desenvolvimento Java 8 e Python 3.9 instalados para executar clientes como aplicações Java, aplicações Python e Beeline.
O endereço IP do seu cliente deve estar na lista de permissões do cluster AnalyticDB for MySQL.
Observações de uso
Se um grupo de recursos Spark Interactive estiver parado, o cluster o reiniciará quando você executar a primeira consulta Spark SQL. A primeira consulta pode entrar em fila durante a inicialização.
O Spark não lê nem grava nos bancos de dados INFORMATION_SCHEMA e MYSQL. Não use esses bancos como banco de dados inicial de conexão.
Verifique se a conta de banco de dados usada para enviar jobs Spark SQL tem acesso ao banco de dados de destino. Caso contrário, a consulta falhará.
Preparação
-
Obtenha o endpoint do grupo de recursos Spark Interactive.
Faça login no console do AnalyticDB for MySQL. No canto superior esquerdo do console, selecione uma região. No painel de navegação à esquerda, clique em Clusters. Localize o cluster que deseja gerenciar e clique no ID do cluster.
No painel de navegação à esquerda, escolha e clique na aba Resource Groups.
-
Localize o grupo de recursos e clique em Details na coluna Actions para visualizar o endpoint interno e o endpoint público. Clique no ícone
ao lado de um endpoint para copiá-lo ou clique no ícone
dentro dos parênteses do número da Port para copiar a string de conexão JDBC.Nos casos abaixo, clique em Apply for Endpoint ao lado de Public Address para solicitar manualmente um endpoint público.
A ferramenta cliente usada para enviar jobs Spark SQL está implantada na sua máquina local ou em um servidor externo.
A ferramenta cliente usada para enviar jobs Spark SQL está implantada em uma instância ECS, e a instância ECS e o cluster AnalyticDB for MySQL não estão na mesma VPC.
As informações de conexão também incluem campos como porta pública e VPC (padrão:
10000), VPC ID, VSwitch ID, classe do driver (org.apache.hive.jdbc.HiveDriver) e URL de download do driver.
Análise interativa
Console
Se você utilizar um HiveMetastore autogerenciado, crie um banco de dados chamado default no AnalyticDB for MySQL e selecione-o como banco de dados ao executar jobs Spark SQL no console.
Faça login no console do AnalyticDB for MySQL. No canto superior esquerdo do console, selecione uma região. No painel de navegação à esquerda, clique em Clusters. Localize o cluster que deseja gerenciar e clique no ID do cluster.
No painel de navegação à esquerda, escolha .
-
Selecione o mecanismo Spark e o grupo de recursos Spark Interactive criado e execute a seguinte instrução Spark SQL:
SHOW DATABASES;
SDK
Ao executar instruções Spark SQL usando um SDK, os resultados da consulta são gravados como arquivos em um bucket OSS especificado. Em seguida, consulte os dados no console do OSS ou baixe os arquivos de resultado para o seu computador. O exemplo a seguir mostra como chamar o SDK em Python.
-
Execute o comando a seguir para instalar o SDK.
pip install alibabacloud-adb20211201 -
Execute os comandos a seguir para instalar as dependências.
pip install oss2 pip install loguru -
Conecte-se ao cluster e execute instruções 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)Parâmetros:
_ak: O AccessKey ID da sua conta Alibaba Cloud ou de um usuário RAM com permissões de acesso ao AnalyticDB for MySQL. Para obter informações sobre como conseguir um AccessKey ID e um AccessKey secret, consulte Contas e permissões.
_sk: O AccessKey secret da sua conta Alibaba Cloud ou de um usuário RAM com permissões de acesso ao AnalyticDB for MySQL. Para obter informações sobre como conseguir um AccessKey ID e um AccessKey secret, consulte Contas e permissões.
_region: O ID da região onde reside o seu cluster AnalyticDB for MySQL.
_db: O ID do cluster AnalyticDB for MySQL.
_rg_name: O nome do grupo de recursos Spark Interactive.
-
oss_location (opcional): O caminho do OSS onde os arquivos de resultado da consulta são armazenados.
Se você não especificar este parâmetro, poderá visualizar apenas as primeiras cinco linhas do resultado da consulta no Log na página .
Aplicações
Hive JDBC
-
No arquivo pom.xml, configure a dependência Maven.
<dependency> <groupId>org.apache.hive</groupId> <artifactId>hive-jdbc</artifactId> <version>2.3.9</version> </dependency> -
Estabeleça uma conexão e execute instruções 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")); } } }Parâmetros:
String de conexão JDBC: A string de conexão JDBC do grupo de recursos Spark Interactive obtida na seção Preparação. Substitua default pelo nome do banco de dados ao qual deseja se conectar.
Nome de usuário: A conta do AnalyticDB for MySQL para o AnalyticDB for MySQL.
Senha: A senha da conta do AnalyticDB for MySQL para o AnalyticDB for MySQL.
PyHive
-
Instale o cliente PyHive.
pip install pyhive -
Estabeleça uma conexão e execute instruções 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())Parâmetros:
Endpoint: O endpoint do grupo de recursos Spark Interactive obtido na seção Preparação.
Porta: A porta do grupo de recursos Spark Interactive, que é 10000.
Nome do grupo de recursos: O nome do grupo de recursos Spark Interactive.
Nome de usuário: A conta do AnalyticDB for MySQL para o AnalyticDB for MySQL.
Senha: A senha da conta do AnalyticDB for MySQL para o AnalyticDB for MySQL.
Clientes
Além dos clientes Beeline, DBeaver, DBVisualizer e DataGrip descritos neste tópico, você também pode realizar análises interativas em ferramentas de agendamento de fluxo de trabalho, como Airflow, Azkaban e DolphinScheduler.
Beeline
-
Conecte-se ao grupo de recursos Spark Interactive.
Use o seguinte formato de comando:
!connect <JDBC-connection-string> <username> <password>String de conexão JDBC: A string de conexão JDBC do grupo de recursos Spark Interactive obtida na seção Preparação. Substitua default pelo nome do banco de dados ao qual deseja se conectar.
Nome de usuário: A conta do AnalyticDB for MySQL para o AnalyticDB for MySQL.
Senha: A senha da conta do AnalyticDB for MySQL para o AnalyticDB for MySQL.
Exemplo:
!connect jdbc:hive2://amv-bp1c3em7b2e****-spark.ads.aliyuncs.com:10000/adb_test spark_resourcegroup/AdbSpark14**** Spark23****Uma conexão bem-sucedida retorna a seguinte saída:
Connected to: Spark SQL (version 3.2.0) Driver: Hive JDBC (version 2.3.9) Transaction isolation: TRANSACTION_REPEATABLE_READ -
Execute uma instrução Spark SQL.
SHOW TABLES;
DBeaver
Abra o cliente DBeaver e escolha .
Na página Connect to a database, selecione Apache Spark e clique em Next.
-
Configure as Hadoop/Apache Spark connection settings conforme descrito abaixo:
Parâmetro
Descrição
Método de conexão
Selecione URL.
JDBC URL
Insira a string de conexão JDBC obtida na seção Preparação.
ImportanteSubstitua
defaultna string de conexão pelo nome do seu banco de dados.Nome de usuário
A conta do AnalyticDB for MySQL para AnalyticDB for MySQL.
Senha
A senha da conta do AnalyticDB for MySQL para AnalyticDB for MySQL.
-
Após configurar os parâmetros, clique em Test Connection.
ImportanteNa primeira vez que você testar a conexão, o DBeaver solicitará o download dos drivers necessários. Clique em Download para baixá-los.
Após o teste de conexão ser bem-sucedido, clique em Finish.
Na aba Database Navigator, expanda a fonte de dados e clique no banco de dados.
-
No editor de código à direita, insira uma instrução SQL e clique no ícone
para executá-la.SHOW TABLES;+-----------+-----------+-------------+ | namespace | tableName | isTemporary | +-----------+-----------+-------------+ | db | test | [] | +-----------+-----------+-------------+
DBVisualizer
Abra o cliente DBVisualizer e escolha .
Na página Driver Manager, selecione Hive e clique no ícone
.-
Na aba Driver Settings, configure os seguintes parâmetros:
Parâmetro
Descrição
Nome
Um nome personalizado para a fonte de dados Hive.
Formato da URL
Insira a string de conexão JDBC obtida na seção Preparação.
ImportanteSubstitua
defaultna string de conexão pelo nome do seu banco de dados.Classe do driver
Selecione org.apache.hive.jdbc.HiveDriver.
ImportanteApós configurar os parâmetros, clique em Start Download para baixar o driver.
Depois que o driver for baixado, escolha .
-
Na caixa de diálogo Create Database Connection from Database URL, configure os parâmetros descritos na tabela a seguir.
Parâmetro
Descrição
Database URL
Insira a string de conexão JDBC obtida na seção Preparação.
ImportanteSubstitua
defaultna string de conexão pelo nome do seu banco de dados.Classe do driver
Selecione a fonte de dados Hive criada na Etapa 3.
-
Na página Connection, configure os seguintes parâmetros de conexão e clique em Connect.
Parâmetro
Descrição
Nome
Por padrão, este parâmetro é definido como o nome da fonte de dados Hive criada na Etapa 3. Você pode personalizar o nome.
Notas
Insira observações.
Tipo de driver
Selecione Hive.
Database URL
Insira a string de conexão JDBC obtida na seção Preparação.
ImportanteSubstitua
defaultna string de conexão pelo nome do seu banco de dados.Database userid
A conta do AnalyticDB for MySQL para AnalyticDB for MySQL.
Database password
A senha da conta do AnalyticDB for MySQL para AnalyticDB for MySQL.
NotaMantenha os demais parâmetros com seus valores padrão.
Após estabelecer a conexão, na aba Database, expanda a fonte de dados e clique no banco de dados.
-
No editor de código à direita, insira uma instrução SQL e clique no ícone
para executá-la.SHOW TABLES;+-----------+-----------+-------------+ | namespace | tableName | isTemporary | +-----------+-----------+-------------+ | db | test | false | +-----------+-----------+-------------+
DataGrip
Abra o cliente DataGrip, escolha e crie um projeto.
-
Adicione uma fonte de dados.
Clique no ícone
e escolha .-
Na caixa de diálogo Data Sources and Drivers exibida, configure os seguintes parâmetros e clique em OK.
Defina Driver como Apache Spark e Authentication como User & Password.
Parâmetro
Descrição
Nome
O nome da fonte de dados, que pode ser personalizado. Este tópico usa
adbtestcomo exemplo.Host
Insira a string de conexão JDBC obtida na seção Preparação.
ImportanteSubstitua
defaultna string de conexão pelo nome do seu banco de dados.Porta
A porta do grupo de recursos Spark Interactive, que é 10000.
Usuário
A conta do AnalyticDB for MySQL para AnalyticDB for MySQL.
Senha
A senha da conta do AnalyticDB for MySQL para AnalyticDB for MySQL.
Schema
O nome do banco de dados no cluster AnalyticDB for MySQL.
-
Execute instruções Spark SQL.
Na lista de fontes de dados, clique com o botão direito na fonte de dados criada na Etapa 2 e escolha .
-
No painel Console exibido, execute uma instrução Spark SQL.
SHOW TABLES;
Ferramentas de BI
É possível realizar análises interativas em ferramentas de BI como Redash, Power BI e Metabase.