O Hologres Spark Connector permite conectar o EMR Serverless Spark ao Hologres mediante as configurações necessárias. Este tópico descreve como ler e gravar dados no Hologres em um ambiente Serverless Spark.
Limitações
O conector do Spark exige a versão 1,3 ou posterior do Hologres. Verifique a versão da sua instância na página Instance Details no console do Hologres. Se a versão for anterior à 1,3, atualize sua instância ou entre no grupo DingTalk do Hologres (ID: 32314975) para solicitar a atualização.
Métodos de acesso
Há duas formas de acessar o Hologres. Escolha o método mais adequado às suas necessidades:
Método de acesso | Descrição | Cenários | Referências |
Método 1: Configuração no nível de tarefa/sessão | Configure as informações de conexão do Hologres (URL JDBC, nome de usuário, senha e outros parâmetros) separadamente em cada tarefa ou sessão. |
| Este tópico |
Método 2: Configuração unificada via catálogos de dados (recomendado) | Adicione um catálogo de dados do Hologres pelo recurso Data Catalogs do EMR Serverless Spark. Após adicionar o catálogo, todos os jobs e sessões no workspace podem acessar os dados autorizados por padrão. Nota Somente a versão de engine esr-4.9.0 e posteriores são compatíveis. |
|
Se o seu workspace exigir acesso frequente e de longo prazo aos dados do Hologres, recomendamos o Método 2 (catálogos de dados) para reduzir configurações repetitivas e aumentar a eficiência do desenvolvimento.
Procedimento
Etapa 1: Obter e carregar o JAR hologres-connector-spark****
As versões esr-4.8.0 e posteriores do EMR Serverless Spark já incluem o conector do Hologres integrado. Esta etapa só é necessária se você usar uma versão anterior à esr-4.8.0.
Para ler e gravar no Hologres, o Spark requer um arquivo JAR de conector. Baixe-o no Repositório Central Maven. Este tópico utiliza a versão 1.5.6: hologres-connector-spark-3.x-1.5.6-jar-with-dependencies.jar.
Carregue o arquivo JAR hologres-connector-spark baixado no OSS. Para obter instruções, consulte Upload simples.
Etapa 2: Adicionar uma conexão de rede
-
Obtenha as informações de rede.
Acesse a página do Hologres e acesse os detalhes da instância de destino do Hologres para encontrar as informações de VPC e vSwitch.
-
Adicione uma conexão de rede.
O Serverless Spark precisa de uma conexão de rede com o cluster do Hologres para acessar o serviço. Para mais informações sobre conexões de rede, consulte Conectividade de rede entre o EMR Serverless Spark e outras VPCs.
Etapa 3: Criar um banco de dados e uma tabela no Hologres
Conecte-se à instância do Hologres. Para detalhes, consulte Conectar-se a uma instância.
-
Na aba SQL Editor, insira as seguintes instruções SQL em uma nova consulta temporária e execute-as.
-- Create a database. CREATE DATABASE testdb; -- Create a table. CREATE TABLE "public"."test" ( "id" text NULL, "name" text NULL); -- Insert data. INSERT INTO public.test VALUES ('1001','jack'),('1002','tony'),('1003','mike'); -- Query data. SELECT * FROM public.test
Etapa 4: Ler e gravar no Hologres
Exemplo 1: Sessão SQL
Este exemplo demonstra como ler e gravar no Hologres usando uma sessão SQL.
-
Crie uma sessão SQL. Para detalhes, consulte Gerenciar sessões SQL.
Ao criar a sessão, selecione a conexão de rede criada na etapa anterior na lista de conexões de rede. Na seção Spark Configuration, adicione os seguintes parâmetros para carregar o hologres-connector-spark.
# Add the hologres-connector JAR file (only required for versions before esr-4.8.0). spark.emr.serverless.user.defined.jars oss://<bucket>/hologres-connector-spark-3.x-<version>.jar # Configure the Hologres catalog. spark.sql.catalog.hologres_external_test_db com.alibaba.hologres.spark3.HoloTableCatalog spark.sql.catalog.hologres_external_test_db.username *** spark.sql.catalog.hologres_external_test_db.password *** spark.sql.catalog.hologres_external_test_db.jdbcurl jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/testdbA tabela a seguir descreve os parâmetros.
Parâmetro
Exemplo
Descrição
spark.emr.serverless.user.defined.jarsoss://<bucket>/hologres-connector-spark-3.x-<version>.jarEspecifica o caminho para o arquivo JAR definido pelo usuário.
spark.sql.catalog.hologres_external_test_dbcom.alibaba.hologres.spark3.HoloTableCatalogNo Spark 3.x, este parâmetro configura uma fonte de dados do Hologres como um catálogo externo. Trata-se de um valor fixo.
spark.sql.catalog.hologres_external_test_db.usernameLTAI******O AccessKey ID da sua conta Alibaba Cloud. Recomendamos usar o gerenciamento de segredos para lidar com informações sensíveis. Para detalhes, consulte Gerenciar informações sensíveis usando o gerenciamento de segredos.
spark.sql.catalog.hologres_external_test_db.passwordmXYV******O AccessKey Secret da sua conta Alibaba Cloud. Recomendamos usar o gerenciamento de segredos para lidar com informações sensíveis. Para detalhes, consulte Gerenciar informações sensíveis usando o gerenciamento de segredos.
spark.sql.catalog.hologres_external_test_db.jdbcurljdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/testdbA URL de conexão JDBC da instância do Hologres.
Você pode personalizar a parte
hologres_external_test_dbnos nomes dos parâmetros. -
Na página Data Development, crie um job SparkSQL e selecione a sessão SQL criada no canto superior direito.
Para mais informações, consulte Desenvolvimento SparkSQL.
-
Copie o código abaixo para a nova aba SparkSQL e clique em Run.
-- Switch to the testdb database. USE hologres_external_test_db; -- Write data. INSERT INTO `public`.test VALUES ('1004','tom'); -- Query data. SELECT * FROM `public`.test;
Exemplo 2: Job de streaming
Este exemplo em PySpark mostra como ler dados do Kafka e gravá-los no Hologres como um job de streaming.
Certifique-se de que a conexão de rede entre o Kafka e o Hologres esteja ativa. Recomendamos implantar o Kafka e o Hologres na mesma VPC e vSwitch.
-
Neste exemplo de código, substitua as informações do Kafka e a tabela do Hologres pelos seus valores reais.
from pyspark.sql import SparkSession from pyspark.sql.functions import col # Configure your Kafka information. servers = "alikafka-serverless-cn-xxxxx-vpc.alikafka.aliyuncs.com:9092" # Replace with your Kafka bootstrap servers. topic = "topic-name" # Replace with your Kafka topic. # Create a SparkSession. spark = SparkSession.builder \ .appName("test read kafka") \ .getOrCreate() # Read the Kafka stream. df = spark \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", servers) \ .option("subscribe", topic) \ .load() # Define a function to write to Hologres (called for each micro-batch). def write_to_hologres(batch_df, batch_id): print(f"Writing batch {batch_id} to Hologres...") batch_df.write \ .format("hologres") \ .mode("append") \ .insertInto("hologres_external_test_db.public.test") # Replace with your Hologres table. # Convert the key and value to strings and write the stream. query = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") \ .writeStream \ .foreachBatch(write_to_hologres) \ .outputMode("append") \ .trigger(processingTime='30 seconds') \ .start() # Wait for the streaming query to terminate (this blocks cell execution in a notebook). query.awaitTermination() -
Carregue o arquivo.
Na página Artifacts, clique em Upload File.
Na caixa de diálogo Upload File, clique na área de upload para selecionar o arquivo Python da etapa anterior ou arraste o arquivo para a área de upload.
-
Crie e execute o job de streaming.
Na página Development, clique no ícone
(Create).Na caixa de diálogo exibida, insira um Name, selecione PySpark na lista Application (Streaming) e clique em OK.
-
Na nova aba de desenvolvimento, configure os parâmetros a seguir e mantenha os demais com as configurações padrão. Em seguida, clique em Publish.
Parâmetro
Descrição
Main Python Resources
Selecione o arquivo Python carregado na etapa anterior.
Engine Version
Selecione uma versão compatível do Spark. Este exemplo usa
esr-4.6.0.Network Connection
Selecione a conexão de rede criada na Etapa 2.
Spark Configuration
# Add the hologres-connector JAR file (only required for versions before esr-4.8.0). spark.emr.serverless.user.defined.jars oss://<bucket>/test_script/hologres-connector-spark-3.x-1.5.6-jar-with-dependencies.jar # Configure the Hologres catalog. spark.sql.catalog.hologres_external_test_db com.alibaba.hologres.spark3.HoloTableCatalog spark.sql.catalog.hologres_external_test_db.username *** spark.sql.catalog.hologres_external_test_db.password *** spark.sql.catalog.hologres_external_test_db.jdbcurl jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/testdbPara uma descrição detalhada dos parâmetros, consulte Exemplo 1: Sessão SQL.
Após publicar o job, clique em Go to O&M. Na página aberta, clique em Start.
-
Verifique o resultado.
-
Envie mensagens para o Kafka.

Consulte os dados usando Spark SQL.

-
Exemplo 3: Sessão de Notebook
-
Crie uma sessão de Notebook. Para detalhes, consulte Gerenciar sessões de Notebook.
Ao criar a sessão, selecione a conexão de rede criada na etapa anterior na lista de conexões de rede. Na seção Spark Configuration, adicione o seguinte parâmetro para carregar o hologres-connector-spark.
# Add the hologres-connector JAR file (only required for versions before esr-4.8.0). spark.emr.serverless.user.defined.jars oss://<bucket>/hologres-connector-spark-3.x-<version>.jar Na página Data Development, crie um job Notebook e selecione a sessão de Notebook criada no canto superior direito.
-
Copie o código abaixo para a nova aba do Notebook e clique em
.import pandas as pd from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, IntegerType, LongType # 1. Prepare a Pandas DataFrame. pdf = pd.DataFrame({ "id": ["1006"], "name": ["sl"] }) # 2. Convert to a PySpark DataFrame. # (Optional: Explicitly define the schema to ensure correct data types) schema = StructType([ StructField("id", StringType(), True), StructField("name", StringType(), True) ]) df = spark.createDataFrame(pdf, schema=schema) # Write to Hologres. df.write \ .format("hologres") \ .option("username", "LTAI******") \ .option("password", "mXYV******") \ .option("jdbcurl", "jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/testdb") \ .option("table", "test") \ .mode("append") \ .save() # Read data. readDf = spark.read\ .format("hologres") \ .option("username", "LTAI******") \ .option("password", "mXYV******") \ .option("jdbcurl", "jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/testdb") \ .option("table", "test") \ .load() readDf.select("id", "name").show(10)A tabela a seguir descreve os parâmetros.
Parâmetro
Exemplo
Descrição
usernameLTAI******O AccessKey ID da sua conta Alibaba Cloud. Recomendamos usar o gerenciamento de segredos para lidar com informações sensíveis. Para detalhes, consulte Gerenciar informações sensíveis usando o gerenciamento de segredos.
passwordmXYV******O AccessKey Secret da sua conta Alibaba Cloud. Recomendamos usar o gerenciamento de segredos para lidar com informações sensíveis. Para detalhes, consulte Gerenciar informações sensíveis usando o gerenciamento de segredos.
jdbcurljdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/testdbA URL de conexão JDBC da instância do Hologres.
Verifique o resultado.

Comandos comuns de catálogo do Hologres
No Serverless Spark, use um catálogo do Hologres para conectar um banco de dados do Hologres ao Spark SQL como um catálogo externo. Cada catálogo está estritamente vinculado a um único banco de dados do Hologres, e não há suporte para acesso entre bancos de dados (não é possível usar o mesmo catálogo para acessar múltiplos bancos de dados do Hologres). A estrutura lógica dentro de um catálogo é consistente com o Hologres:
|
Conceito Spark |
Conceito Hologres |
Descrição |
|
catalog |
database |
Por exemplo, |
|
namespace |
schema |
Por exemplo, |
|
table |
table |
Especifique explicitamente |
Uso de um catálogo do Hologres
Um catálogo do Hologres no Spark mapeia exatamente para um banco de dados do Hologres e não pode ser alterado após a criação.
USE hologres_external_test_db;
Listagem de todos os namespaces
No Spark, um namespace corresponde a um schema no Hologres. O schema padrão é public. Use o comando USE para alterar o schema padrão de uma sessão.
-- View all namespaces in the Hologres catalog, which correspond to all schemas in Hologres.
SHOW NAMESPACES;
Listagem de tabelas em um namespace
-
Liste todas as tabelas.
SHOW TABLES; -
Liste as tabelas em um namespace específico.
USE test_schema; SHOW TABLES; -- Or use: SHOW TABLES IN test_schema;
Documentos relacionados
Para mais informações sobre como usar o Spark para ler e gravar no Hologres, consulte Ler e gravar dados no Hologres usando Spark.