Este tópico descreve como sincronizar um banco de dados MySQL inteiro com o Kafka. Essa abordagem reduz a carga que vários jobs exercem sobre o banco de dados MySQL.
Contexto
Uma tabela source do MySQL CDC captura dados do MySQL e sincroniza alterações em tempo real. Esse cenário é comum em computação complexa, como quando uma tabela serve como tabela de dimensão em uma operação JOIN com outras tabelas. Uma única tabela MySQL pode ser dependência de vários jobs. Quando múltiplos jobs processam dados da mesma tabela MySQL, o banco de dados abre várias conexões, o que gera pressão significativa no servidor MySQL e na rede.
Como funciona
Para reduzir a pressão no banco de dados MySQL upstream, o Realtime Compute for Apache Flink pode sincronizar um banco de dados MySQL inteiro com o Kafka. Essa solução introduz o Kafka como camada intermediária e utiliza um job de ingestão de dados do Flink CDC para sincronizar os dados.
Em um único job, os dados de um banco de dados MySQL upstream são sincronizados com o Kafka em tempo real. Cada tabela MySQL é gravada em um tópico Kafka correspondente no modo upsert. Os jobs downstream usam então um conector Kafka upsert para ler dados dos tópicos, em vez de acessar as tabelas MySQL diretamente. Esse método reduz eficazmente a pressão que múltiplos jobs exercem sobre o banco de dados MySQL.

Limitações
Cada tabela MySQL sincronizada deve ter uma chave primária.
É possível utilizar clusters Kafka autogerenciados, clusters EMR Kafka ou ApsaraMQ for Kafka. Se você usar o ApsaraMQ for Kafka, conecte-se apenas por meio de um endpoint padrão.
O espaço de armazenamento do cluster Kafka deve ser maior que o das tabelas source. Caso contrário, pode haver perda de dados devido a armazenamento insuficiente. Os tópicos criados para sincronização de banco de dados são tópicos compactados. Em um tópico compactado, apenas a mensagem mais recente para cada chave de mensagem é retida, mas os dados nunca expiram. Isso significa que o tópico compactado armazena um volume de dados aproximadamente equivalente ao tamanho da tabela source.
Cenário de exemplo
Por exemplo, em um cenário de análise de revisão de pedidos em tempo real, suponha que existam três tabelas: uma tabela de usuários (user), uma tabela de pedidos (order) e uma tabela de feedback de usuários (feedback). As tabelas contêm os dados mostrados na figura a seguir.
Para exibir informações de pedidos e avaliações de usuários, junte a tabela user para recuperar nomes de usuários do campo name. O exemplo SQL a seguir demonstra essa operação.
-- Join order information with the user table to display the username and product name for each order.
SELECT order.id as order_id, product, user.name as user_name
FROM order LEFT JOIN user
ON order.user_id = user.id;
-- Join reviews with the user table to display the content of each review and the corresponding username.
SELECT feedback.id as feedback_id, comment, user.name as user_name
FROM feedback LEFT JOIN user
ON feedback.user_id = user.id;
Ambos os jobs SQL anteriores utilizam a tabela user. Durante a execução, ambos os jobs leem os dados completos e incrementais do MySQL. Uma leitura completa exige a criação de uma conexão MySQL, e uma leitura incremental exige a criação de um cliente Binlog. À medida que o número de jobs aumenta, a demanda por recursos de conexão MySQL e clientes Binlog também cresce, exercendo pressão significativa no banco de dados upstream. Para aliviar essa pressão, utilize um job de ingestão de dados do Flink CDC para sincronizar dados do banco de dados MySQL upstream com o Kafka em tempo real, permitindo o consumo por múltiplos jobs downstream.
Pré-requisitos
O Realtime Compute for Apache Flink está ativado. Para obter mais informações, consulte Ativar o Realtime Compute for Apache Flink.
O ApsaraMQ for Kafka está ativado. Para obter mais informações, consulte Implantar uma instância do ApsaraMQ for Kafka.
O ApsaraDB RDS for MySQL está ativado. Para obter mais informações, consulte Criar uma instância do ApsaraDB RDS for MySQL.
Os serviços Realtime Compute for Apache Flink, ApsaraDB RDS for MySQL e ApsaraMQ for Kafka estão na mesma VPC. Se estiverem em VPCs diferentes, habilite o acesso à rede entre VPCs ou use endpoints públicos. Para obter mais informações, consulte Como acesso outros serviços entre VPCs? e Como acesso a Internet?.
Se você acessar recursos como usuário RAM ou usando uma função RAM, deverá ter as permissões necessárias.
Preparação
Preparar a fonte de dados MySQL
-
Crie um banco de dados ApsaraDB RDS for MySQL. Para obter mais informações, consulte Criar um banco de dados.
Crie um banco de dados chamado
order_dwpara a instância de destino. -
Prepare a fonte de dados MySQL CDC.
Na página de detalhes da instância, clique em Log on to Database na parte superior da página.
Na caixa de diálogo de login do DMS exibida, insira o nome de usuário e a senha da conta do banco de dados criada e clique em Login.
Após fazer login, clique duas vezes no banco de dados
order_dwno painel esquerdo para trocar de banco de dados.-
No SQL Console, insira as instruções DDL para criar as três tabelas de negócios e as instruções para inserir dados.
CREATE TABLE `user` ( id bigint not null primary key, name varchar(50) not null ); CREATE TABLE `order` ( id bigint not null primary key, product varchar(50) not null, user_id bigint not null ); CREATE TABLE `feedback` ( id bigint not null primary key, user_id bigint not null, comment varchar(50) not null ); -- Prepare data INSERT INTO `user` VALUES(1, 'Tom'),(2, 'Jerry'); INSERT INTO `order` VALUES (1, 'Football', 2), (2, 'Basket', 1); INSERT INTO `feedback` VALUES (1, 1, 'Good.'), (2, 2, 'Very good');
Clique em Execute e, em seguida, clique em Execute.
Procedimento
-
Crie e inicie um job de ingestão de dados do Flink CDC para sincronizar dados do banco de dados MySQL upstream com o Kafka em tempo real, permitindo o consumo por múltiplos jobs downstream. O job de sincronização de banco de dados cria tópicos automaticamente. Defina os nomes dos tópicos usando o módulo
route. Os tópicos usam as configurações padrão do cluster Kafka para o número de partições e réplicas, ecleanup.policyé definido comocompact.Nomes de tópicos padrão
Por padrão, os tópicos Kafka criados pelo job de sincronização de banco de dados usam o formato de nomenclatura
{database_name}.{table_name}. O job a seguir cria três tópicos:order_dw.user,order_dw.ordereorder_dw.feedback.-
Na página , crie um job de ingestão de dados do Flink CDC e copie o código a seguir no editor YAML.
source: type: mysql name: MySQL Source hostname: #{hostname} port: 3306 username: #{usernmae} password: #{password} tables: order_dw.\.* server-id: 28601-28604 # (Optional) Synchronize data from newly created tables during the incremental phase. scan.binlog.newly-added-table.enabled: true # (Optional) Synchronize table and field comments. include-comments.enabled: true # (Optional) Prioritize dispatching unbounded splits to avoid potential TaskManager OutOfMemory issues. scan.incremental.snapshot.unbounded-chunk-first.enabled: true # (Optional) Enable parsing filters to speed up reads. scan.only.deserialize.captured.tables.changelog.enabled: true sink: type: upsert-kafka name: upsert-kafka Sink properties.bootstrap.servers: xxxx.alikafka.aliyuncs.com:9092 # The following parameters are required for ApsaraMQ for Kafka. aliyun.kafka.accessKeyId: #{ak} aliyun.kafka.accessKeySecret: #{sk} aliyun.kafka.instanceId: #{instanceId} aliyun.kafka.endpoint: #{endpoint} aliyun.kafka.regionId: #{regionId} No canto superior direito, clique em Deploy para implantar o job.
No painel de navegação à esquerda, clique em . Na coluna Actions do job alvo, clique em Start, selecione Initial Mode e clique em Start.
Nomes de tópicos por tabela
Utilize o módulo
routepara especificar um nome de tópico para cada tabela. O job a seguir cria três tópicos:user1,order2efeedback3.-
Na página , crie um job de ingestão de dados do Flink CDC e copie o código a seguir no editor YAML.
source: type: mysql name: MySQL Source hostname: #{hostname} port: 3306 username: #{usernmae} password: #{password} tables: order_dw.\.* server-id: 28601-28604 # (Optional) Synchronize data from newly created tables during the incremental phase. scan.binlog.newly-added-table.enabled: true # (Optional) Synchronize table and field comments. include-comments.enabled: true # (Optional) Prioritize dispatching unbounded splits to avoid potential TaskManager OutOfMemory issues. scan.incremental.snapshot.unbounded-chunk-first.enabled: true # (Optional) Enable parsing filters to speed up reads. scan.only.deserialize.captured.tables.changelog.enabled: true route: - source-table: order_dw.user sink-table: user1 - source-table: order_dw.order sink-table: order2 - source-table: order_dw.feedback sink-table: feedback3 sink: type: upsert-kafka name: upsert-kafka Sink properties.bootstrap.servers: xxxx.alikafka.aliyuncs.com:9092 # The following parameters are required for ApsaraMQ for Kafka. aliyun.kafka.accessKeyId: #{ak} aliyun.kafka.accessKeySecret: #{sk} aliyun.kafka.instanceId: #{instanceId} aliyun.kafka.endpoint: #{endpoint} aliyun.kafka.regionId: #{regionId} No canto superior direito, clique em Deploy para implantar o job.
No painel de navegação à esquerda, selecione , clique em Start na coluna Actions do job alvo, selecione Initial Mode e clique em Start.
Nomes de tópicos em lote
Use o módulo
routepara especificar um padrão para os nomes dos tópicos gerados. O job a seguir cria três tópicos:topic_user,topic_orderetopic_feedback.-
Na página , crie um job de ingestão de dados do Flink CDC e copie o código a seguir no editor YAML.
source: type: mysql name: MySQL Source hostname: #{hostname} port: 3306 username: #{usernmae} password: #{password} tables: order_dw.\.* server-id: 28601-28604 # (Optional) Synchronize data from newly created tables during the incremental phase. scan.binlog.newly-added-table.enabled: true # (Optional) Synchronize table and field comments. include-comments.enabled: true # (Optional) Prioritize dispatching unbounded splits to avoid potential TaskManager OutOfMemory issues. scan.incremental.snapshot.unbounded-chunk-first.enabled: true # (Optional) Enable parsing filters to speed up reads. scan.only.deserialize.captured.tables.changelog.enabled: true route: - source-table: order_dw.\.* sink-table: topic_<> replace-symbol: <> sink: type: upsert-kafka name: upsert-kafka Sink properties.bootstrap.servers: xxxx.alikafka.aliyuncs.com:9092 # The following parameters are required for ApsaraMQ for Kafka. aliyun.kafka.accessKeyId: #{ak} aliyun.kafka.accessKeySecret: #{sk} aliyun.kafka.instanceId: #{instanceId} aliyun.kafka.endpoint: #{endpoint} aliyun.kafka.regionId: #{regionId} No canto superior direito, clique em Deploy para implantar o job.
No painel de navegação à esquerda, clique em . Clique em Start na coluna Actions do job alvo, selecione Initial Mode e clique em Start.
-
-
Consuma dados do Kafka em tempo real.
O job de ingestão de dados grava dados do banco de dados MySQL upstream no Kafka em formato JSON. Múltiplos jobs downstream podem então consumir dados de um único tópico para recuperar o estado mais recente das tabelas do banco de dados. Consuma os dados das tabelas sincronizadas com o Kafka de uma das seguintes maneiras:
Por catálogo
Leia dados de um tópico Kafka usando-o como uma tabela source.
-
Na página , crie um job SQL de streaming e copie o código a seguir no editor SQL.
CREATE TEMPORARY TABLE print_user_proudct( order_id BIGINT, product STRING, user_name STRING ) WITH ( 'connector'='print', 'logger'='true' ); CREATE TEMPORARY TABLE print_user_feedback( feedback_id BIGINT, `comment` STRING, user_name STRING ) WITH ( 'connector'='print', 'logger'='true' ); BEGIN STATEMENT SET; -- Required when writing to multiple sinks. -- Join order information with the user table in the Kafka JSON Catalog to display the username and product name for each order. INSERT INTO print_user_proudct SELECT `order`.key_id as order_id, value_product as product, `user`.value_name as user_name FROM `kafka-catalog`.`kafka`.`order`/*+OPTIONS('properties.group.id'='<yourGroupName>', 'scan.startup.mode'='earliest-offset')*/ as `order` -- Specify the group and startup mode. LEFT JOIN `kafka-catalog`.`kafka`.`user`/*+OPTIONS('properties.group.id'='<yourGroupName>', 'scan.startup.mode'='earliest-offset')*/ as `user` -- Specify the group and startup mode. ON `order`.value_user_id = `user`.key_id; -- Join reviews with the user table to display the content of each review and the corresponding username. INSERT INTO print_user_feedback SELECT feedback.key_id as feedback_id, value_comment as `comment`, `user`.value_name as user_name FROM `kafka-catalog`.`kafka`.feedback/*+OPTIONS('properties.group.id'='<yourGroupName>', 'scan.startup.mode'='earliest-offset')*/ as feedback -- Specify the group and startup mode. LEFT JOIN `kafka-catalog`.`kafka`.`user`/*+OPTIONS('properties.group.id'='<yourGroupName>', 'scan.startup.mode'='earliest-offset')*/ as `user` -- Specify the group and startup mode. ON feedback.value_user_id = `user`.key_id; END; -- Required when writing to multiple sinks.Este exemplo usa o conector Print para imprimir os resultados diretamente. Também é possível enviar os resultados para uma tabela de resultados que use outro conector para análise posterior. Para obter mais informações sobre a sintaxe para gravar em múltiplos sinks, consulte Instrução INSERT INTO.
NotaAo usar este método diretamente, podem ocorrer alterações de schema. Como resultado, o schema analisado pelo catálogo JSON do Kafka pode diferir do schema da tabela MySQL correspondente. Por exemplo, campos excluídos ainda podem aparecer e alguns campos podem ter valores nulos.
O schema lido do catálogo consiste em campos dos dados consumidos. Se um campo for excluído, mas suas mensagens não tiverem expirado, o campo ainda poderá aparecer com um valor nulo. Nenhum tratamento especial é necessário para esse caso.
No canto superior direito, clique em Deploy para implantar o job.
No painel de navegação à esquerda, clique em , clique em Start na coluna Actions do job alvo, selecione Initial Mode e clique em Start.
Por tabela temporária
Defina um schema personalizado e leia dados de uma tabela temporária.
-
Na página , crie um job SQL de streaming e copie o código a seguir no editor SQL.
CREATE TEMPORARY TABLE user_source ( key_id BIGINT, value_name STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'user', 'properties.bootstrap.servers' = '<yourKafkaBrokers>', 'scan.startup.mode' = 'earliest-offset', 'key.format' = 'json', 'value.format' = 'json', 'key.fields' = 'key_id', 'key.fields-prefix' = 'key_', 'value.fields-prefix' = 'value_', 'value.fields-include' = 'EXCEPT_KEY', 'value.json.infer-schema.flatten-nested-columns.enable' = 'false', 'value.json.infer-schema.primitive-as-string' = 'false' ); CREATE TEMPORARY TABLE order_source ( key_id BIGINT, value_product STRING, value_user_id BIGINT ) WITH ( 'connector' = 'kafka', 'topic' = 'order', 'properties.bootstrap.servers' = '<yourKafkaBrokers>', 'scan.startup.mode' = 'earliest-offset', 'key.format' = 'json', 'value.format' = 'json', 'key.fields' = 'key_id', 'key.fields-prefix' = 'key_', 'value.fields-prefix' = 'value_', 'value.fields-include' = 'EXCEPT_KEY', 'value.json.infer-schema.flatten-nested-columns.enable' = 'false', 'value.json.infer-schema.primitive-as-string' = 'false' ); CREATE TEMPORARY TABLE feedback_source ( key_id BIGINT, value_user_id BIGINT, value_comment STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'feedback', 'properties.bootstrap.servers' = '<yourKafkaBrokers>', 'scan.startup.mode' = 'earliest-offset', 'key.format' = 'json', 'value.format' = 'json', 'key.fields' = 'key_id', 'key.fields-prefix' = 'key_', 'value.fields-prefix' = 'value_', 'value.fields-include' = 'EXCEPT_KEY', 'value.json.infer-schema.flatten-nested-columns.enable' = 'false', 'value.json.infer-schema.primitive-as-string' = 'false' ); CREATE TEMPORARY TABLE print_user_proudct( order_id BIGINT, product STRING, user_name STRING ) WITH ( 'connector'='print', 'logger'='true' ); CREATE TEMPORARY TABLE print_user_feedback( feedback_id BIGINT, `comment` STRING, user_name STRING ) WITH ( 'connector'='print', 'logger'='true' ); BEGIN STATEMENT SET; -- Required when writing to multiple sinks. -- Join order information with the user table from the Kafka JSON Catalog to display the username and product name for each order. INSERT INTO print_user_proudct SELECT order_source.key_id as order_id, value_product as product, user_source.value_name as user_name FROM order_source LEFT JOIN user_source ON order_source.value_user_id = user_source.key_id; -- Join reviews with the user table to display the content of each review and the corresponding username. INSERT INTO print_user_feedback SELECT feedback_source.key_id as feedback_id, value_comment as `comment`, user_source.value_name as user_name FROM feedback_source LEFT JOIN user_source ON feedback_source.value_user_id = user_source.key_id; END; -- Required when writing to multiple sinks.Este exemplo usa o conector Print para imprimir os resultados diretamente. Também é possível enviar os resultados para uma tabela de resultados de outro conector para análise posterior. Para obter mais informações sobre a sintaxe para gravar em múltiplos sinks, consulte Instrução INSERT INTO.
A tabela a seguir descreve os parâmetros de configuração para a tabela temporária.
Parâmetro
Descrição
Observações
connectorO tipo de conector.
Defina o valor como
kafka.topicO nome do tópico correspondente.
Deve ser consistente com a descrição do catálogo JSON do Kafka.
properties.bootstrap.serversOs endereços dos brokers Kafka.
O formato é
host:port,host:port,host:port, separado por vírgulas (,).scan.startup.modeA posição inicial para leitura de dados do Kafka.
Valores válidos:
-
earliest-offset: Inicia a leitura a partir do offset mais antigo disponível. -
latest-offset: Inicia a leitura a partir do offset mais recente. -
group-offsets(padrão): Lê a partir do offset confirmado para o grupo especificado por properties.group.id. -
timestamp: Lê a partir do timestamp especificado por scan.startup.timestamp-millis.
-
specific-offsets: Inicia a leitura a partir dos offsets especificados em scan.startup.specific-offsets.
Nota
Este parâmetro entra em vigor quando o job inicia sem um estado salvo. Quando o job reinicia ou se recupera de um checkpoint, ele prioriza a leitura a partir do estado salvo.
key.formatO formato que o conector Flink Kafka usa para serializar ou desserializar a chave da mensagem Kafka.
Defina o valor como
json.key.fieldsOs campos na tabela source ou de resultados que correspondem à chave da mensagem Kafka.
Use ponto e vírgula (;) para separar vários nomes de campos. Por exemplo,
field1;field2.key.fields-prefixUm prefixo personalizado para todos os campos de chave de mensagem Kafka para evitar conflitos de nomenclatura com campos de valor de mensagem ou campos de metadados.
Este valor deve ser consistente com o valor do parâmetro key.fields-prefix do Catálogo JSON do Kafka.
value.formatO formato que o conector Flink Kafka usa para serializar ou desserializar o valor da mensagem Kafka.
Defina o valor como
json.value.fields-prefixUm prefixo personalizado para todos os campos de valor de mensagem Kafka para evitar conflitos de nomenclatura com campos de chave de mensagem ou campos de metadados.
Deve corresponder ao valor do parâmetro value.fields-prefix do Catálogo JSON do Kafka.
value.fields-includeA política para lidar com campos de chave de mensagem no valor da mensagem.
Defina o valor como
EXCEPT_KEY. Isso indica que o valor da mensagem não inclui os campos da chave da mensagem.value.json.infer-schema.flatten-nested-columns.enableEspecifica se deve expandir recursivamente colunas JSON aninhadas no valor da mensagem Kafka.
O valor do parâmetro infer-schema.flatten-nested-columns.enable do Catálogo correspondente.
value.json.infer-schema.primitive-as-stringEspecifica se deve inferir todos os tipos primitivos como String no valor da mensagem Kafka.
O valor do parâmetro infer-schema.primitive-as-string do Catálogo correspondente.
-
No canto superior direito, clique em Deploy para implantar o job.
No painel de navegação à esquerda, clique em , clique em Start na coluna Actions do job alvo, selecione Initial Mode e clique em Start.
-
-
Visualize os resultados do job.
No painel de navegação à esquerda, clique em e clique no job alvo.
Na aba Job Log, na aba Running TaskManagers, clique na tarefa que possui o Path, ID que você deseja visualizar.
-
Clique em Logs e pesquise informações de log relacionadas a
PrintSinkOutputWriter.Pesquise nos logs a saída de
PrintSinkOutputWriter. A saída contém quatro registros de dados unidos:+I[1, Good., Tom],+I[2, Very good, Jerry],+I[2, Basket, Tom]e+I[1, Football, Jerry]. Isso indica que as junções entre a tabela de usuários e as tabelas de pedidos e feedback foram bem-sucedidas.