Todos os produtos
Search
Central de documentação

AnalyticDB:Sincronizar dados do Kafka com o recurso de sincronização de dados do APS (Recomendado)

Última atualização: Jun 27, 2026

O recurso de sincronização de dados do AnalyticDB Pipeline Service (APS) permite ingerir mensagens do ApsaraMQ for Kafka no AnalyticDB for MySQL em tempo quase real, a partir de qualquer offset escolhido. Esse recurso oferece saída de dados em tempo quase real, arquivamento completo de dados históricos e análises elásticas. Após o início de uma tarefa de sincronização, os dados são confirmados (committed) a cada 5 minutos por padrão. Assim, o primeiro lote de dados ingeridos fica disponível para consulta após aproximadamente 5 minutos.

O sistema aceita apenas mensagens do Kafka formatadas em json.

Limitações

  • O sistema aceita apenas mensagens do Kafka em formato json. Outros formatos causam erros de sincronização.

  • Alterações no schema da tabela do Kafka não se propagam automaticamente para o AnalyticDB for MySQL. Aplique as mudanças de DDL manualmente.

  • A API do Kafka trunca dados de amostra maiores que 8 KB, o que impede a análise da amostra e a geração automática de mapeamentos de campos.

  • Os dados gravados no AnalyticDB for MySQL só ficam visíveis após a execução de uma operação de Commit. Como o intervalo padrão de Commit é de 5 minutos, aguarde pelo menos esse período após iniciar a tarefa antes de consultar o primeiro lote de dados.

  • Caminhos do OSS em tarefas de sincronização diferentes não podem compartilhar um prefixo. Por exemplo, oss://testBucketName/test/sls1/ e oss://testBucketName/test/ entram em conflito e causam a sobrescrita de dados.

  • Se uma tarefa de sincronização falhar e os dados do tópico do Kafka tiverem expirado, não será possível recuperar as informações removidas ao reiniciar a tarefa. Para reduzir esse risco, aumente o período de retenção de dados do tópico. Em caso de falha na tarefa, entre em contato com o suporte técnico imediatamente.

Pré-requisitos

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

  • Um cluster do AnalyticDB for MySQL nas edições Enterprise, Basic ou Data Lakehouse

  • Um grupo de recursos de job

  • Uma conta de banco de dados para o cluster (consulte a tabela abaixo)

  • Uma instância do ApsaraMQ for Kafka na mesma região do cluster do AnalyticDB for MySQL

  • Um tópico do Kafka com mensagens já enviadas (consulte Início rápido do ApsaraMQ for Kafka)

Requisitos de conta de banco de dados por tipo de conta:

Tipo de conta

Contas necessárias

Etapas adicionais

Conta Alibaba Cloud

Apenas Conta privilegiada

Nenhuma

Usuário do Resource Access Management (RAM)

Conta privilegiada + conta padrão

Associar a conta padrão ao usuário RAM

Faturamento

O uso do recurso de sincronização de dados gera as seguintes taxas:

Como funciona

  1. Adicione uma source de dados Kafka para identificar a instância e o tópico do ApsaraMQ for Kafka de onde os dados serão lidos.

  2. Crie um link de dados (job de sincronização) que mapeie as mensagens do Kafka para uma tabela do AnalyticDB for MySQL e configure a análise de json, as chaves de partição e o offset do consumidor.

  3. Inicie a tarefa. Ela começa a consumir dados do Kafka a partir do offset selecionado e os grava na tabela de destino.

  4. Consulte os dados ingeridos usando Spark SQL.

Etapa 1: Criar uma source de dados

Ignore esta etapa se você já adicionou uma source de dados Kafka . Vá diretamente para a Etapa 2: Criar um link de dados .
  1. Faça login no console do AnalyticDB for MySQL. No canto superior esquerdo, selecione uma região. No painel de navegação à esquerda, clique em Clusters e, em seguida, clique no ID do cluster.

  2. No painel de navegação à esquerda, escolha Data Ingestion > Data Sources.

  3. Clique em Create Data Source.

  4. Na página Create Data Source, configure os seguintes parâmetros:

    Parâmetro

    Descrição

    Data Source Type

    Selecione Kafka.

    Data Source Name

    Gerado automaticamente com base no tipo de source e na hora atual. Altere conforme necessário.

    Data Source Description

    Insira uma descrição, como o cenário de negócios ou o escopo dos dados.

    Deployment Mode

    Apenas Alibaba Cloud Instance é compatível.

    Kafka Instance

    O ID da instância Kafka. Encontre-o na página Instances do console do ApsaraMQ for Kafka.

    Kafka Topic

    O nome do tópico. Localize-o na página Topics da instância de destino no console do ApsaraMQ for Kafka.

    Message Data Format

    Apenas json é compatível.

  5. Clique em Create.

Etapa 2: Criar um link de dados

  1. No painel de navegação à esquerda, clique em Simple Log Service/Kafka Data Synchronization.

  2. Clique em Create Synchronization Job.

  3. Na página Create Synchronization Job, configure as três seções abaixo.

Configurações de source e destino

Parâmetro

Descrição

Job Name

Gerado automaticamente com base no tipo de source e na hora atual. Modifique conforme necessário.

Data Source

Selecione uma source de dados Kafka existente ou crie uma nova.

Destination Type

Escolha onde armazenar os dados sincronizados. Consulte Escolher um tipo de destino abaixo.

ADB Lake Storage

O armazenamento de lake para dados do AnalyticDB for MySQL. Selecione na lista suspensa ou clique em Automatically Created para criar um. Obrigatório quando o Destination Type for Data Lake - AnalyticDB Lake Storage.

OSS Path

O caminho de armazenamento no OSS para dados do lake do AnalyticDB for MySQL. Obrigatório quando o Destination Type for Data Lake - User OSS. Selecione uma pasta vazia — o caminho não pode ser alterado após a criação e não deve compartilhar prefixo com o caminho do OSS de nenhuma outra tarefa de sincronização.

Storage Format

PAIMON (disponível apenas quando o Destination Type for Data Lake - User OSS) ou ICEBERG.

Escolher um tipo de destino

Tipo de destino

Quando usar

Observações

Data Lake - AnalyticDB Lake Storage (recomendado)

Configurações padrão em que você deseja que o AnalyticDB for MySQL gerencie o armazenamento

Ative o recurso de armazenamento de lake primeiro. Apenas o formato ICEBERG é compatível.

Data Lake - User OSS

Quando você precisa usar seu próprio bucket do OSS

Os formatos PAIMON e ICEBERG são compatíveis.

Configurações de banco de dados e tabela de destino

Parâmetro

Descrição

Database Name

O banco de dados de destino no AnalyticDB for MySQL. Um novo banco de dados será criado se não existir nenhum com esse nome. Caso já exista, os dados serão sincronizados nele. Consulte Limites para convenções de nomenclatura. Se o Storage Format for PAIMON, o banco de dados deve atender aos requisitos listados abaixo.

Table Name

A tabela de destino no AnalyticDB for MySQL. Uma nova tabela será criada se não existir nenhuma com esse nome. Se já houver uma tabela com o mesmo nome, a tarefa de sincronização falhará. Consulte Limites para convenções de nomenclatura.

Sample Data

Os dados mais recentes recuperados automaticamente do tópico do Kafka. Devem estar no formato json.

Parsed JSON Layers

O número de níveis aninhados de json a serem analisados. Valores válidos: 0 (sem análise), 1 (padrão), 2, 3, 4. Consulte Níveis de análise de json e inferência de schema.

Schema Field Mapping

O schema inferido a partir dos dados de amostra após a análise. Ajuste os nomes e tipos dos campos de destino, ou adicione e remova campos conforme necessário.

Partition Key Settings

(Opcional) Uma chave de partição para a tabela de destino. Particione por tempo de log ou lógica de negócios para melhorar o desempenho de ingestão e consulta. Se deixado em branco, a tabela não terá partições.

Requisitos de banco de dados PAIMON

Se o Storage Format for PAIMON, o banco de dados de destino deve atender a todas as condições a seguir. Caso contrário, a tarefa de sincronização falhará.

  • Deve ser um banco de dados externo criado com CREATE EXTERNAL DATABASE <database_name>.

  • O parâmetro DBPROPERTIES deve incluir catalog = paimon.

  • O parâmetro DBPROPERTIES deve incluir adb.paimon.warehouse. Exemplo: adb.paimon.warehouse=oss://testBucketName/aps/data.

  • O parâmetro DBPROPERTIES deve incluir LOCATION com o sufixo .db no nome do banco de dados. Exemplo: LOCATION=oss://testBucketName/aps/data/kafka_paimon_external_db.db/. O diretório do bucket neste caminho já deve existir, e o caminho deve incluir .db após o nome do banco de dados — caso contrário, as consultas XIHE falharão.

Configurações de sincronização

Parâmetro

Descrição

Starting Consumer Offset for Incremental Synchronization

O ponto no Kafka a partir do qual a tarefa começa a consumir dados. Earliest offset (begin_cursor): consome a partir dos dados disponíveis mais antigos. Latest offset (end_cursor): consome apenas os dados mais recentes. Custom offset: consome a partir da primeira mensagem no horário específico selecionado ou posterior a ele.

Job Resource Group

O grupo de recursos de job onde a tarefa será executada.

ACUs for Incremental Synchronization

A quantidade de ACUs alocadas do grupo de recursos de job. Mínimo: 2 ACUs. Máximo: as ACUs restantes disponíveis no grupo de recursos. Uma contagem maior de ACUs melhora o desempenho de ingestão e a estabilidade da tarefa.

Advanced Settings

Configurações personalizadas. Entre em contato com o suporte técnico para ativar.

Exemplo de dedução de ACU: Se um grupo de recursos de job tem um máximo de 48 ACUs e uma tarefa existente já utiliza 8, a nova tarefa poderá usar no máximo 40 ACUs.

  1. Clique em Submit.

Etapa 3: Iniciar a tarefa de sincronização de dados

  1. Na página Simple Log Service/Kafka Data Synchronization, localize a tarefa criada e clique em Start na coluna Actions.

  2. Clique em Search. A tarefa foi iniciada com sucesso quando seu status mudar para Running.

O primeiro lote de dados estará disponível para consulta pelo menos 5 minutos após o início da tarefa, pois os dados são confirmados em intervalos de 5 minutos por padrão.

Etapa 4: Analisar os dados

Após a sincronização dos dados, use Spark SQL para consultá-los no AnalyticDB for MySQL. Para mais informações, consulte Editor de desenvolvimento Spark e Desenvolvimento de aplicações Spark offline.

  1. No painel de navegação à esquerda, escolha Job Development > Spark JAR Development.

  2. Insira suas instruções Spark SQL no modelo padrão e clique em Run Now. Veja um exemplo abaixo:

    -- Example of Spark SQL. Modify the content and run your Spark program.
    
    conf spark.driver.resourceSpec=medium;
    conf spark.executor.instances=2;
    conf spark.executor.resourceSpec=medium;
    conf spark.app.name=Spark SQL Test;
    conf spark.adb.connectors=oss;
    
    -- SQL statements
    show tables from lakehouse20220413156_adbTest;
  3. (Opcional) Na aba Applications, clique em Logs na coluna Actions para visualizar os logs de execução do Spark SQL.

Etapa 5 (Opcional): Gerenciar a source de dados

Acesse Data Ingestion > Data Sources. As seguintes operações estão disponíveis na coluna Actions:

Operação

Descrição

Create Job

Cria uma tarefa de sincronização ou migração de dados para esta source.

View

Visualiza a configuração da source de dados.

Edit

Edita o nome e a descrição da source de dados.

Delete

Exclui a source de dados. Se houver alguma tarefa de sincronização ou migração associada a esta source, exclua essa tarefa primeiro na página Simple Log Service/Kafka Data Synchronization.

Níveis de análise de json e inferência de schema

A configuração Parsed JSON Layers controla quantos níveis aninhados de uma mensagem json são expandidos em campos de destino separados.

Mensagem de exemplo:

{
  "name" : "zhangle",
  "age" : 18,
  "device" : {
    "os" : {
        "test":lag,
        "member":{
             "fa":zhangsan,
             "mo":limei
           }
         },
    "brand" : "none",
    "version" : "11.4.2"
  }
}
Importante

Pontos (.) nos nomes dos campos são substituídos automaticamente por underscores (_) nos nomes dos campos de destino.

Nível 0 — Sem análise. A mensagem json inteira é enviada como um único campo.

Campo JSON

Valor

Nome do campo de destino

__value__

{"name":"zhangle","age":18,"device":{...}}

__value__

Nível 1 (padrão) — Os campos de nível superior são expandidos.

Campo JSON

Valor

Nome do campo de destino

name

zhangle

name

age

18

age

device

{"os":{...},"brand":"none","version":"11.4.2"}

device

Nível 2 — Dois níveis expandidos. Campos não aninhados são enviados diretamente; campos aninhados se expandem para seus subcampos.

Campo JSON

Valor

Nome do campo de destino

name

zhangle

name

age

18

age

device.os

{"test":"lag","member":{...}}

device_os

device.brand

none

device_brand

device.version

11.4.2

device_version

Nível 3

Campo JSON

Valor

Nome do campo de destino

name

zhangle

name

age

18

age

device.os.test

lag

device_os_test

device.os.member

{"fa":"zhangsan","mo":"limei"}

device_os_member

device.brand

none

device_brand

device.version

11.4.2

device_version

Nível 4

Campo JSON

Valor

Nome do campo de destino

name

zhangle

name

age

18

age

device.os.test

lag

device_os_test

device.os.member.fa

zhangsan

device_os_member_fa

device.os.member.mo

lime

device_os_member_mo

device.brand

none

device_brand

device.version

11.4.2

device_version

Próximos passos