Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Flink CDC: Sincronizar banco de dados MySQL inteiro para o Kafka

Última atualização: Jun 27, 2026

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.

图片 1

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.mysql database

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

Preparação

Preparar a fonte de dados MySQL

  1. 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_dw para a instância de destino.

  2. Prepare a fonte de dados MySQL CDC.

    1. Na página de detalhes da instância, clique em Log on to Database na parte superior da página.

    2. 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.

    3. Após fazer login, clique duas vezes no banco de dados order_dw no painel esquerdo para trocar de banco de dados.

    4. 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');
  3. Clique em Execute e, em seguida, clique em Execute.

Procedimento

  1. 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, e cleanup.policy é definido como compact.

    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.order e order_dw.feedback.

    1. Na página Development > Data Ingestion, 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}
    2. No canto superior direito, clique em Deploy para implantar o job.

    3. No painel de navegação à esquerda, clique em Operations Center > Deployments. 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 route para especificar um nome de tópico para cada tabela. O job a seguir cria três tópicos: user1, order2 e feedback3.

    1. Na página Development > Data Ingestion, 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}
    2. No canto superior direito, clique em Deploy para implantar o job.

    3. No painel de navegação à esquerda, selecione Operations Center > Deployments, 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 route para especificar um padrão para os nomes dos tópicos gerados. O job a seguir cria três tópicos: topic_user, topic_order e topic_feedback.

    1. Na página Development > Data Ingestion, 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}
    2. No canto superior direito, clique em Deploy para implantar o job.

    3. No painel de navegação à esquerda, clique em Operations Center > Deployments. Clique em Start na coluna Actions do job alvo, selecione Initial Mode e clique em Start.

  1. 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.

    1. Na página Development > ETL, 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.

      Nota

      Ao 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.

    2. No canto superior direito, clique em Deploy para implantar o job.

    3. No painel de navegação à esquerda, clique em O&M Center > Deployments, 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.

    1. Na página Development > ETL, 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

      connector

      O tipo de conector.

      Defina o valor como kafka.

      topic

      O nome do tópico correspondente.

      Deve ser consistente com a descrição do catálogo JSON do Kafka.

      properties.bootstrap.servers

      Os endereços dos brokers Kafka.

      O formato é host:port,host:port,host:port, separado por vírgulas (,).

      scan.startup.mode

      A 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.format

      O formato que o conector Flink Kafka usa para serializar ou desserializar a chave da mensagem Kafka.

      Defina o valor como json.

      key.fields

      Os 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-prefix

      Um 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.format

      O formato que o conector Flink Kafka usa para serializar ou desserializar o valor da mensagem Kafka.

      Defina o valor como json.

      value.fields-prefix

      Um 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-include

      A 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.enable

      Especifica 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-string

      Especifica 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.

    2. No canto superior direito, clique em Deploy para implantar o job.

    3. No painel de navegação à esquerda, clique em Operations Center > Deployments, clique em Start na coluna Actions do job alvo, selecione Initial Mode e clique em Start.

  2. Visualize os resultados do job.

    1. No painel de navegação à esquerda, clique em Operations Center > Deployments e clique no job alvo.

    2. Na aba Job Log, na aba Running TaskManagers, clique na tarefa que possui o Path, ID que você deseja visualizar.

    3. 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.

Documentos relacionados