Todos os produtos
Search
Central de documentação

DataHub:Canal plug-in

Última atualização: Jun 28, 2026

O Canal lê os logs binários (binlogs) do MySQL para capturar alterações incrementais e transmiti-las para tópicos do DataHub usando o protocolo Kafka.

Como o Canal funciona com o DataHub

Nota

O Canal suporta as versões 5.1.x, 5.5.x, 5.6.x, 5.7.x e 8.0.x do MySQL como bancos de dados source.

O DataHub é compatível com o protocolo Kafka, permitindo que o Canal grave dados incrementais do MySQL diretamente em tópicos do DataHub. O pacote do Canal personalizado para o DataHub inclui duas alterações em relação à versão open-source:

  • Remoção da lógica que substituía pontos (.) por underscores (_) nos nomes dos tópicos Kafka. Isso permite que o Canal mapeie o TopicName do Kafka para o tópico correto do DataHub usando o formato ProjectName.TopicName.

  • Adição da variável de ambiente -Djava.security.auth.login.config=$kafka_jaas_conf ao script de inicialização para oferecer suporte à autenticação PLAIN Simple Authentication and Security Layer (SASL).

Pré-requisitos

Antes de começar, verifique se você possui:

  • Um banco de dados MySQL executando a versão 5.1.x, 5.5.x, 5.6.x, 5.7.x ou 8.0.x. Consulte QuickStart para obter detalhes sobre a configuração.

  • Um tópico do DataHub do tipo TUPLE criado no projeto de destino. Para requisitos de tópicos, consulte Compatibilidade com Kafka.

Para todos os parâmetros de configuração do Canal e opções avançadas, consulte Canal.

Configure o Canal para transmitir dados do MySQL para o DataHub

Etapa 1: Baixe o pacote do Canal

Baixe o pacote do Canal personalizado para o DataHub: canal.deployer-1.1.5-SNAPSHOT.tar.gz

Use este pacote em vez da versão padrão open-source do Canal. Versões do Canal não modificadas para o DataHub podem não conseguir gravar dados no serviço.

Etapa 2: Extraia o pacote

mkdir -p /usr/local/canal
tar -zxvf canal.deployer-1.1.5-SNAPSHOT.tar.gz -C /usr/local/canal

Etapa 3: Configure o Canal

Edite três arquivos de configuração antes de iniciar o Canal.

**3,1 Configuração da instância: conf/example/instance.properties**

# Modify the database information as needed.
#################################################
...
canal.instance.master.address=192.168.1.20:3306
# username/password: the username and password of the database.
...
canal.instance.dbUsername = canal
canal.instance.dbPassword = canal
...
# mq config
canal.mq.topic=test_project.test_topic
# Specify a dynamic topic based on the database name or table name.
#canal.mq.dynamicTopic=mytest,.*,mytest.user,mytest\\..*,.*\\..*
canal.mq.partition=0
# hash partition config
#canal.mq.partitionsNum=3
# Database name.Table name: the unique primary key. Multiple tables are separated with comma (,).
#canal.mq.partitionHash=mytest.person:id,mytest.role:id
#################################################

Defina canal.mq.topic como <ProjectName>.<TopicName> — o formato separado por ponto mapeia diretamente para um tópico do DataHub. Para roteamento dinâmico de tópicos com base em nomes de banco de dados ou tabelas e para configurar chaves de hash de partição, consulte Parâmetros relacionados a MQ.

**3,2 Configuração do Canal: conf/canal.properties**

# ...
canal.serverMode = kafka
# ...
kafka.bootstrap.servers = dh-cn-hangzhou.aliyuncs.com:9092
kafka.acks = all
kafka.compression.type = none
kafka.batch.size = 16384
kafka.linger.ms = 1
kafka.max.request.size = 1048576
kafka.buffer.memory = 33554432
kafka.max.in.flight.requests.per.connection = 1
kafka.retries = 0

kafka.security.protocol = SASL_SSL
kafka.sasl.mechanism = PLAIN

Os seguintes parâmetros são obrigatórios:

Parâmetro

Valor

Descrição

canal.serverMode

kafka

Instrui o Canal a usar o protocolo producer do Kafka

kafka.bootstrap.servers

Endpoint do DataHub

O endpoint do DataHub na região onde o tópico de destino reside. Para endpoints disponíveis, consulte Compatibilidade com Kafka.

kafka.security.protocol

SASL_SSL

Necessário para autenticação no DataHub

kafka.sasl.mechanism

PLAIN

Obrigatório para autenticação SASL/PLAIN do DataHub

Todos os outros parâmetros são opcionais e podem ser ajustados para atender aos seus requisitos de throughput.

**3,3 Configuração do Java Authentication and Authorization Service (JAAS): conf/kafka_client_producer_jaas.conf**

kafkaClient {
  org.apache.kafka.common.security.plain.PlainLoginModule required
  username="accessId"
  password="accessKey";
};

Etapa 4: Inicie o Canal e verifique o fluxo de dados

Iniciar o Canal

cd /usr/local/canal/
sh bin/startup.sh

Verificar se o Canal está em execução

Consulte o log principal do Canal para confirmar que o servidor foi iniciado com sucesso:

2013-02-05 22:45:27.967 [main] INFO  com.alibaba.otter.canal.deployer.CanalLauncher - ## start the canal server.
2013-02-05 22:45:28.113 [main] INFO  com.alibaba.otter.canal.deployer.CanalController - ## start the canal server[10.1.29.120:11111]
2013-02-05 22:45:28.210 [main] INFO  com.alibaba.otter.canal.deployer.CanalLauncher - ## the canal server is running now ......

Execute vi logs/canal/canal.log para visualizar o log principal. Execute vi logs/example/example.log para verificar o log da instância:

2013-02-05 22:50:45.636 [main] INFO  c.a.o.c.i.spring.support.PropertyPlaceholderConfigurer - Loading properties file from class path resource [canal.properties]
2013-02-05 22:50:45.641 [main] INFO  c.a.o.c.i.spring.support.PropertyPlaceholderConfigurer - Loading properties file from class path resource [example/instance.properties]
2013-02-05 22:50:45.803 [main] INFO  c.a.otter.canal.instance.spring.CanalInstanceWithSpring - start CannalInstance for 1-example 
2013-02-05 22:50:45.810 [main] INFO  c.a.otter.canal.instance.spring.CanalInstanceWithSpring - start successful....

Verificar a captura de dados

Execute vi logs/example/meta.log para confirmar que o Canal está capturando alterações no banco de dados. Um registro é gerado no arquivo meta.log para cada inserção, exclusão e modificação no banco de dados. Visualize o arquivo meta.log para checar se o Canal coletou os dados.

tail -f example/meta.log
2020-07-29 09:21:05.110 - clientId:1001 cursor:[log.000001,29723,1591190230000,1,] address[/127.0.0.1:3306]
2020-07-29 09:23:46.109 - clientId:1001 cursor:[log.000001,30047,1595985825000,1,] address[localhost/127.0.0.1:3306]
2020-07-29 09:24:50.547 - clientId:1001 cursor:[log.000001,30047,1595985825000,1,] address[/127.0.0.1:3306]
2020-07-29 09:26:45.547 - clientId:1001 cursor:[log.000001,30143,1595986005000,1,] address[localhost/127.0.0.1:3306]
2020-07-29 09:30:04.546 - clientId:1001 cursor:[log.000001,30467,1595986204000,1,] address[localhost/127.0.0.1:3306]
2020-07-29 09:30:16.546 - clientId:1001 cursor:[log.000001,30734,1595986215000,1,] address[localhost/127.0.0.1:3306]
2020-07-29 09:30:36.547 - clientId:1001 cursor:[log.000001,31001,1595986236000,1,] address[localhost/127.0.0.1:3306]

Parar o Canal

cd /usr/local/canal/
sh bin/stop.sh

Exemplo

O exemplo a seguir mostra como alterações de linha no MySQL aparecem como registros do DataHub após o processamento pelo Canal.

Schema do tópico do DataHub

O tópico de destino é do tipo TUPLE com o seguinte schema:

+-------+------+----------+-------------+
| Index | name |   type   |  allow NULL |
+-------+------+----------+-------------+
|   0   |  key |  STRING  |     true    |
|   1   |  val |  STRING  |     true    |
+-------+------+----------+-------------+

Tabela source do MySQL

mysql> desc orders;
+-------+---------+------+-----+---------+-------+
| Field | Type    | Null | Key | Default | Extra |
+-------+---------+------+-----+---------+-------+
| oid   | int(11) | YES  |     | NULL    |       |
| pid   | int(11) | YES  |     | NULL    |       |
| num   | int(11) | YES  |     | NULL    |       |
+-------+---------+------+-----+---------+-------+
3 rows in set (0.00 sec)

Registro resultante no DataHub

Após a execução de mysql> insert into orders values(1,2,3);, o Canal grava a alteração no DataHub. O campo key é nulo e o campo val contém o evento de alteração completo como uma string JSON:

{
    "data":[
        {
            "oid":"1",
            "pid":"2",
            "num":"3"
        }
    ],
    "database":"ggtt",
    "es":1591092305000,
    "id":2,
    "isDdl":false,
    "mysqlType":{
        "oid":"int(11)",
        "pid":"int(11)",
        "num":"int(11)"
    },
    "old":null,
    "pkNames":null,
    "sql":"",
    "sqlType":{
        "oid":4,
        "pid":4,
        "num":4
    },
    "table":"orders",
    "ts":1591092305813,
    "type":"INSERT"
}

Próximos passos

Com o Canal em execução e dados aparecendo no meta.log, acesse o console do DataHub e verifique o tópico de destino para confirmar a chegada dos registros. Em seguida, conecte um consumidor — como um job do Flink ou uma aplicação downstream — para processar os dados incrementais.