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
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 |
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: |
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:
O padrão é |
sink.properties.* | Parâmetros de importação para o Stream Load. Exemplo: Para especificar um separador de colunas, use 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 Atualmente, o ApsaraDB for SelectDB só garante que dados importados via copy into sejam |
enable.2pc | Define se o commit de duas fases deve ser ativado para garantir a semântica exactly-once. |
Para outras configurações comuns de sink do Kafka Connect, consulte Configuring Connectors.
Exemplos
Pré-requisitos
-
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 Baixe o doris-kafka-connector-1.0.0.jar e coloque o arquivo JAR no diretório KAFKA_HOME/libs.
Crie uma instância do ApsaraDB for SelectDB. Para mais informações, consulte Create an instance.
Conecte-se a uma instância do ApsaraDB for SelectDB usando o protocolo MySQL. Para mais informações, consulte Connect to an instance.
-
Crie um banco de dados de teste e uma tabela de teste.
-
Crie um banco de dados de teste.
CREATE DATABASE test_db; -
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
-
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 -
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.
-
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 -
Descompacte o arquivo baixado.
tar -zxvf debezium-connector-mysql-1.9.8.Final-plugin.tar.gz Coloque todos os arquivos JAR extraídos no diretório KAFKA_HOME/libs.
-
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=rewriteApós a configuração, o formato padrão do nome do tópico do Kafka será
SERVER_NAME.DATABASE_NAME.TABLE_NAME.NotaPara configurações do Debezium, consulte Debezium connector for MySQL.
-
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=1NotaAo sincronizar dados com o ApsaraDB for SelectDB, crie previamente um banco de dados e uma tabela.
-
Inicie o Kafka Connect.
bin/connect-standalone.sh -daemon config/connect-standalone.properties config/mysql-source.properties config/selectdb-sink.propertiesNotaApó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.