Todos os produtos
Search
Central de documentação

ApsaraDB for SelectDB:Importar dados usando o Kafka

Última atualização: Aug 21, 2026

ApsaraDB for SelectDB oferece suporte ao uso do Doris Kafka Connector para assinar e sincronizar dados do Kafka automaticamente. Este tópico descreve como usar o Doris Kafka Connector para sincronizar dados com o ApsaraDB for SelectDB.

Informações básicas

O Kafka Connect é uma ferramenta para transmitir dados de forma confiável entre o Apache Kafka e outros sistemas. Defina conectores para mover grandes conjuntos de dados para dentro ou para fora do Kafka.

O conector Kafka fornecido pela comunidade Doris executa em um cluster do Kafka Connect. Ele lê dados de um tópico do Kafka e os grava no ApsaraDB for SelectDB.

Em cenários de negócios, os usuários geralmente utilizam o Debezium Connector para enviar dados de alterações de banco de dados ao Kafka ou chamam uma API para gravar dados formatados em JSON no Kafka em tempo real. O Doris Kafka Connector assina automaticamente os dados no Kafka e sincroniza essas informações com o ApsaraDB for SelectDB.

Modos de execução do Kafka Connect

O Kafka Connect possui dois modos de execução.

Modo standalone

Aviso

O modo standalone não é recomendado para ambientes de produção.

Configurar o modo standalone

Configure o arquivo connect-standalone.properties.

# Modify the broker address
bootstrap.servers=127.0.0.1:9092

No diretório de configuração do Kafka, crie um arquivo connect-selectdb-sink.properties e adicione o seguinte conteúdo:

name=test-selectdb-sink
connector.class=org.apache.doris.kafka.connector.DorisSinkConnector
topics=topic_test
doris.topic2table.map=topic_test:test_kafka_tbl
buffer.count.records=10000
buffer.flush.time=120
buffer.size.bytes=5000000
doris.urls=selectdb-cn-4xl3jv1****-public.selectdbfe.rds.aliyuncs.com
doris.http.port=8030
doris.query.port=9030
doris.user=admin
doris.password=****
doris.database=test_db
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter

Iniciar no modo standalone

$KAFKA_HOME/bin/connect-standalone.sh -daemon $KAFKA_HOME/config/connect-standalone.properties $KAFKA_HOME/config/connect-selectdb-sink.properties

Modo distribuído

Configurar o modo distribuído

Configure o arquivo connect-distributed.properties.

# Modify the broker address
bootstrap.servers=127.0.0.1:9092
 
# Modify group.id. The ID must be the same for all workers in the same cluster.
group.id=connect-cluster

Iniciar no modo distribuído

$KAFKA_HOME/bin/connect-distributed.sh -daemon $KAFKA_HOME/config/connect-distributed.properties

Adicionar um conector

curl -i http://127.0.0.1:8083/connectors -H "Content-Type: application/json" -X POST -d '{
  "name":"test-selectdb-sink-cluster",
  "config":{
    "connector.class":"org.apache.doris.kafka.connector.DorisSinkConnector",
    "topics":"topic_test",
    "doris.topic2table.map": "topic_test:test_kafka_tbl",
    "buffer.count.records":"10000",
    "buffer.flush.time":"120",
    "buffer.size.bytes":"5000000",
    "doris.urls":"selectdb-cn-4xl3jv1****-public.selectdbfe.rds.aliyuncs.com",
    "doris.user":"admin",
    "doris.password":"***",
    "doris.database":"test_db",
    "doris.http.port":"8030",
    "doris.query.port":"9030",
    "key.converter":"org.apache.kafka.connect.storage.StringConverter",
    "value.converter":"org.apache.kafka.connect.json.JsonConverter"
  }
}'

Parâmetros

Parâmetro

Descrição

name

Nome do conector. Deve ser uma string sem caracteres de controle ISO e exclusiva no ambiente do Kafka Connect.

connector.class

Nome da classe ou alias do conector. Defina este valor como org.apache.doris.kafka.connector.DorisSinkConnector.

topics

Lista de tópicos de origem separados por vírgula.

doris.topic2table.map

Mapeamento entre tópicos e tabelas. Separe vários mapeamentos com vírgula (,). Exemplo: topic1:tb1,topic2:tb2. Se este parâmetro não for especificado, o conector assumirá que os nomes do tópico e da tabela são iguais.

buffer.count.records

Quantidade de registros armazenados em buffer na memória para cada partição do Kafka antes de serem liberados para o ApsaraDB for SelectDB. O padrão é 10.000 registros.

buffer.flush.time

Intervalo em segundos para liberação do buffer na memória. O padrão é 120.

buffer.size.bytes

Tamanho acumulado em bytes dos registros a serem armazenados em buffer na memória para cada partição do Kafka. O padrão é 5.000.000.

doris.urls

Endpoint de conexão do ApsaraDB for SelectDB.

Obtenha os parâmetros relevantes na página Instance Details > Network Information no console do ApsaraDB for SelectDB.

Exemplo: selectdb-cn-4xl3jv1****-public.selectdbfe.rds.aliyuncs.com

doris.http.port

A porta HTTP padrão do ApsaraDB for SelectDB é 8080.

doris.query.port

A porta do protocolo MySQL do ApsaraDB for SelectDB é 9030 por padrão.

doris.user

Nome de usuário do ApsaraDB for SelectDB.

doris.password

Senha do ApsaraDB for SelectDB.

doris.database

Banco de dados do ApsaraDB for SelectDB onde os dados serão gravados.

key.converter

Classe conversora JSON para a chave.

value.converter

Classe conversora JSON para o valor.

jmx

Define se as métricas internas do conector devem ser recuperadas via JMX. Para mais informações, consulte Doris-Connector-JMX. O padrão é true.

enable.delete

Define se as operações de exclusão devem ser sincronizadas. O padrão é false.

label.prefix

Prefixo do rótulo para dados importados via Stream Load. O padrão é o nome da aplicação do conector.

auto.redirect

Quando ativado, o conector redireciona solicitações de Stream Load através do frontend (FE) para um backend (BE) alvo, eliminando a necessidade de recuperar informações do BE.

load.model

Método de importação de dados. Os seguintes métodos têm suporte:

  • stream_load: Importa dados diretamente para o SelectDB.

  • copy_into: Importa dados para o armazenamento de objetos e depois os carrega no SelectDB.

O padrão é stream_load.

sink.properties.*

Parâmetros de importação para o Stream Load.

Exemplo: Para especificar um separador de colunas, use sink.properties.column_separator=,.

Para mais informações, consulte Stream Load.

delivery.guarantee

Especifica a garantia de consistência de dados ao consumir dados do Kafka e importá-los para o ApsaraDB for SelectDB. Os níveis com suporte são at_least_once e exactly_once, sendo o padrão at_least_once.

Atualmente, o ApsaraDB for SelectDB só garante que dados importados via copy into sejam exactly_once.

enable.2pc

Define se o commit de duas fases deve ser ativado para garantir a semântica exactly-once.

Nota

Para outras configurações comuns de sink do Kafka Connect, consulte Configuring Connectors.

Exemplos

Pré-requisitos

  1. Instale um cluster Apache Kafka ou Confluent Cloud, versão 2.4.0 ou posterior. Este exemplo utiliza um ambiente Kafka de nó único.

    # Download and decompress the package
    wget https://archive.apache.org/dist/kafka/2.4.0/kafka_2.12-2.4.0.tgz
    tar -zxvf kafka_2.12-2.4.0.tgz
    cd kafka_2.12-2.4.0/
    bin/zookeeper-server-start.sh -daemon config/zookeeper.properties 
    bin/kafka-server-start.sh -daemon config/server.properties
  2. Baixe o doris-kafka-connector-1.0.0.jar e coloque o arquivo JAR no diretório KAFKA_HOME/libs.

  3. Crie uma instância do ApsaraDB for SelectDB. Para mais informações, consulte Create an instance.

  4. Conecte-se a uma instância do ApsaraDB for SelectDB usando o protocolo MySQL. Para mais informações, consulte Connect to an instance.

  5. Crie um banco de dados de teste e uma tabela de teste.

    1. Crie um banco de dados de teste.

      CREATE DATABASE test_db;
    2. Crie uma tabela de teste.

      USE test_db;
      CREATE TABLE employees (
          emp_no       int NOT NULL,
          birth_date   date,
          first_name   varchar(20),
          last_name    varchar(20),
          gender       char(2),
          hire_date    date
      )
      UNIQUE KEY(`emp_no`)
      DISTRIBUTED BY HASH(`emp_no`) BUCKETS 1;

Exemplo 1: Sincronizar dados JSON

  1. Configurar o sink do SelectDB

    Tomando o modo standalone como exemplo, crie um arquivo selectdb-sink.properties no diretório de configuração do Kafka e adicione o seguinte conteúdo:

    name=selectdb_sink
    connector.class=org.apache.doris.kafka.connector.DorisSinkConnector
    topics=test_topic
    doris.topic2table.map=test_topic:example_tbl
    buffer.count.records=10000
    buffer.flush.time=120
    buffer.size.bytes=5000000
    doris.urls=selectdb-cn-4xl3jv1****-public.selectdbfe.rds.aliyuncs.com
    doris.http.port=8030
    doris.query.port=9030
    doris.user=admin
    doris.password=***
    doris.database=test_db
    key.converter=org.apache.kafka.connect.storage.StringConverter
    value.converter=org.apache.kafka.connect.json.JsonConverter
    
    #Optional: Configure a dead-letter queue
    errors.tolerance=all
    errors.deadletterqueue.topic.name=test_error
    errors.deadletterqueue.context.headers.enable = true
    errors.deadletterqueue.topic.replication.factor=1
  2. Iniciar o Kafka Connect

    bin/connect-standalone.sh -daemon config/connect-standalone.properties config/selectdb-sink.properties

Exemplo 2: Usar o Debezium para sincronizar dados do MySQL com o ApsaraDB for SelectDB

Em muitos cenários de negócios, é necessário sincronizar dados de um banco de dados operacional em tempo real. Isso requer o uso do mecanismo de captura de dados de alterações (CDC) do banco de dados.

O Debezium é uma ferramenta de CDC baseada no Kafka Connect que se conecta a vários bancos de dados, como MySQL, PostgreSQL, SQL Server, Oracle e MongoDB. Ele envia continuamente alterações de dados para um tópico do Kafka em um formato unificado para consumo em tempo real pelos sinks downstream. Este exemplo utiliza o MySQL.

  1. Baixe o Debezium.

    wget https://repo1.maven.org/maven2/io/debezium/debezium-connector-mysql/1.9.8.Final/debezium-connector-mysql-1.9.8.Final-plugin.tar.gz
  2. Descompacte o arquivo baixado.

    tar -zxvf debezium-connector-mysql-1.9.8.Final-plugin.tar.gz
  3. Coloque todos os arquivos JAR extraídos no diretório KAFKA_HOME/libs.

  4. Configure a origem MySQL.

    Crie um arquivo mysql-source.properties no diretório de configuração do Kafka e adicione o seguinte conteúdo:

    name=mysql-source
    connector.class=io.debezium.connector.mysql.MySqlConnector
    database.hostname=rm-bp17372257wkz****.rwlb.rds.aliyuncs.com
    database.port=3306
    database.user=testuser
    database.password=****
    database.server.id=1
    # A unique identifier for this client in Kafka
    database.server.name=test123
    # The databases and tables to sync. By default, all databases and tables are synced.
    database.include.list=test
    table.include.list=test.test_table
    database.history.kafka.bootstrap.servers=localhost:9092
    # The Kafka topic used to store database schema changes
    database.history.kafka.topic=dbhistory
    transforms=unwrap
    # See https://debezium.io/documentation/reference/stable/transformations/event-flattening.html
    transforms.unwrap.type=io.debezium.transforms.ExtractNewRecordState
    # Record delete events
    transforms.unwrap.delete.handling.mode=rewrite

    Após a configuração, o formato padrão do nome do tópico do Kafka será SERVER_NAME.DATABASE_NAME.TABLE_NAME.

    Nota

    Para configurações do Debezium, consulte Debezium connector for MySQL.

  5. Configure o sink do ApsaraDB for SelectDB.

    Crie um arquivo selectdb-sink.properties no diretório de configuração do Kafka e adicione o seguinte conteúdo:

    name=selectdb-sink
    connector.class=org.apache.doris.kafka.connector.DorisSinkConnector
    topics=test123.test.test_table
    doris.topic2table.map=test123.test.test_table:test_table
    buffer.count.records=10000
    buffer.flush.time=120
    buffer.size.bytes=5000000
    doris.urls=selectdb-cn-4xl3jv1****-public.selectdbfe.rds.aliyuncs.com
    doris.http.port=8030
    doris.query.port=9030
    doris.user=admin
    doris.password=****
    doris.database=test
    key.converter=org.apache.kafka.connect.json.JsonConverter
    value.converter=org.apache.kafka.connect.json.JsonConverter
    
    #Optional: Configure a dead-letter queue
    #errors.tolerance=all
    #errors.deadletterqueue.topic.name=test_error
    #errors.deadletterqueue.context.headers.enable = true
    #errors.deadletterqueue.topic.replication.factor=1
    Nota

    Ao sincronizar dados com o ApsaraDB for SelectDB, crie previamente um banco de dados e uma tabela.

  6. Inicie o Kafka Connect.

    bin/connect-standalone.sh -daemon config/connect-standalone.properties config/mysql-source.properties config/selectdb-sink.properties
    Nota

    Após a inicialização, verifique o arquivo logs/connect.log para confirmar que o service foi iniciado com sucesso.

Uso avançado

Operações do conector

# Check the connector status
curl -i http://127.0.0.1:8083/connectors/test-selectdb-sink-cluster/status -X GET
# Delete the current connector
curl -i http://127.0.0.1:8083/connectors/test-selectdb-sink-cluster -X DELETE
# Pause the current connector
curl -i http://127.0.0.1:8083/connectors/test-selectdb-sink-cluster/pause -X PUT
# Resume the current connector
curl -i http://127.0.0.1:8083/connectors/test-selectdb-sink-cluster/resume -X PUT
# Restart tasks within the connector
curl -i http://127.0.0.1:8083/connectors/test-selectdb-sink-cluster/tasks/0/restart -X POST

Para mais informações, consulte Connect REST Interface.

Fila de mensagens mortas

Por padrão, um erro de conversão causa falha no conector. No entanto, configure o conector para tolerar tais erros ignorando-os. Também é possível gravar os detalhes do erro, a operação com falha e o registro problemático em uma fila de mensagens mortas para análise posterior.

errors.tolerance=all
errors.deadletterqueue.topic.name=test_error_topic
errors.deadletterqueue.context.headers.enable=true
errors.deadletterqueue.topic.replication.factor=1

Para mais informações, consulte Error Reporting in Connect.

Conectar a um cluster Kafka com SSL ativado

Para acessar um cluster Kafka com SSL ativado via Kafka Connect, forneça um arquivo de certificado (client.truststore.jks) para autenticar a chave pública do broker Kafka. Adicione a seguinte configuração ao seu arquivo connect-distributed.properties:

# Connect worker
security.protocol=SSL
ssl.truststore.location=/var/ssl/private/client.truststore.jks
ssl.truststore.password=test1234
 
# Embedded consumer for sink connectors
consumer.security.protocol=SSL
consumer.ssl.truststore.location=/var/ssl/private/client.truststore.jks
consumer.ssl.truststore.password=test1234

Para mais informações sobre como configurar o Kafka Connect para se conectar a um cluster Kafka com SSL ativado, consulte Configure Kafka Connect.