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
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_confao 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 |
|
|
|
Instrui o Canal a usar o protocolo producer do Kafka |
|
|
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. |
|
|
|
Necessário para autenticação no DataHub |
|
|
|
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.