Todos os produtos
Search
Central de documentação

DataWorks:Sincronização em lote de uma única tabela do Kafka para o MaxCompute

Última atualização: Jun 27, 2026

Este tutorial demonstra como configurar uma tarefa recorrente de sincronização offline que lê registros com janela de tempo de um tópico do Kafka e os grava em uma tabela particionada do MaxCompute. A tarefa pode ser executada em qualquer frequência: por minuto, hora, dia, semana ou mês.

Pré-requisitos

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

Limitações

  • A sincronização de dados de origem para tabelas externas do MaxCompute não é suportada.

  • A versão do Kafka deve ser 0.10.2 ou posterior, e 2.2.x ou anterior.

  • O Kafka deve ter timestamps de registro ativados, e cada registro deve conter um timestamp de negócio correto.

Risco de perda de dados

Registros com timestamps anteriores ou iguais à hora de início podem chegar ao tópico do Kafka após o início de uma instância recorrente. Esses registros podem não ser lidos. Se houver atraso nas gravações no tópico do Kafka ou se os timestamps estiverem fora de ordem, poderá ocorrer perda de dados na tarefa de sincronização offline.

Configurar a tarefa de sincronização

Nota

Este tutorial utiliza a nova interface do DataStudio.

Etapa 1: Crie um nó

Siga o guia de Configuração de UI sem código para criar e configurar um nó de sincronização offline. Este tutorial foca nos detalhes de configuração específicos da sincronização do Kafka para o MaxCompute.

Etapa 2: Configure a fonte e o destino dos dados

Configure a fonte de dados (Kafka)

Este tutorial demonstra uma sincronização offline de tabela única do Kafka para o MaxCompute. Para obter uma referência completa de todas as opções de configuração do Kafka Reader, consulte o documento do Kafka Reader.

A tabela a seguir aborda os principais itens de configuração para este tutorial.

Item de configuração

Configuração principal

Topic

Selecione o tópico do Kafka a ser sincronizado. Em um workspace do DataWorks no modo padrão, deve existir um tópico com o mesmo nome nos clusters do Kafka tanto para os ambientes de desenvolvimento quanto de produção.

Nota

Se o tópico estiver ausente no ambiente de desenvolvimento, ele não aparecerá na lista suspensa Topic. Se estiver ausente no ambiente de produção, o agendamento recorrente falhará após a publicação da tarefa, pois ela não conseguirá localizar o tópico para sincronização.

Consumer Group ID

Insira um ID exclusivo para o cluster do Kafka. Esse ID é utilizado para estatísticas e monitoramento no lado do Kafka.

Read Start Offset e Start Time

Defina Read Start Offset como Specific Time e, em seguida, defina Start Time como ${startTime}. Isso define onde a sincronização começa — os registros a partir de ${startTime} são incluídos.

Read End Offset e End Time

Defina Read End Offset como Specific Time e, em seguida, defina End Time como ${endTime}. Os registros até (mas não incluindo) ${endTime} são sincronizados. Tanto ${startTime} quanto ${endTime} são substituídos por timestamps concretos em tempo de execução, com base nas expressões de parâmetro de agendamento configuradas.

Time Zone

Deixe em branco para usar o fuso horário padrão do servidor para a região do DataWorks. Caso tenha alterado o fuso horário de agendamento com o suporte da Alibaba Cloud, selecione esse fuso horário aqui.

Key Type, Value Type, Encoding

Selecione com base nos registros reais presentes no tópico do Kafka.

Synchronization Completion Policy

Controla quando a tarefa de sincronização para de ler. Escolha com base nos seus padrões de tráfego do Kafka — veja a comparação abaixo.

Advanced Configuration

Mantenha os padrões.

Escolhendo uma política de conclusão de sincronização

As duas opções comportam-se de maneira diferente e adequam-se a cenários distintos:

No new data for 1 minute

Stop at the specified end position

Como funciona

A tarefa para quando nenhum novo registro chega em todas as partições durante 1 minuto consecutivo.

A tarefa para assim que atinge o offset correspondente a ${endTime}.

Quando usar

Quando todas as condições a seguir forem verdadeiras: (1) algumas ou todas as partições ficam rotineiramente inativas por longos períodos, como mais de 10 minutos, e (2) nenhum registro com timestamp anterior a ${endTime} será gravado no tópico após o início de cada instância recorrente.

Quando não for possível garantir as condições (1) ou (2) acima. Esta é a opção padrão mais segura.

Risco se mal utilizado

Perda de dados se registros atrasados ainda estiverem sendo gravados quando o silêncio de 1 minuto ocorrer.

Nenhum — a tarefa para de forma confiável na posição final configurada.

Configure o destino dos dados (MaxCompute)

Item de configuração

Configuração principal

Data Source

Exibe a fonte de dados do MaxCompute selecionada na etapa anterior. Em um workspace no modo padrão, os nomes dos projetos de desenvolvimento e produção são exibidos.

Table

Selecione a tabela de destino do MaxCompute. Em um workspace no modo padrão, deve existir uma tabela com o mesmo nome e esquema em ambos os ambientes. Alternativamente, clique em Generate Target Table Schema para permitir que o sistema crie uma tabela automaticamente e, em seguida, ajuste a instrução CREATE TABLE conforme necessário.

Nota

Se a tabela estiver ausente no ambiente de desenvolvimento, ela não aparecerá na lista suspensa. Se estiver ausente no ambiente de produção, a tarefa falhará em tempo de execução. Se os esquemas diferirem entre os ambientes, o mapeamento de colunas poderá gerar gravações de dados incorretas.

Partition

Para tabelas particionadas, insira o valor da chave de partição. Use um valor estático como ds=20220101 ou um parâmetro de agendamento como ds=${partition}, que é substituído automaticamente em tempo de execução.

Etapa 3: Configure o mapeamento de campos

Após selecionar a origem e o destino, mapeie as colunas entre o leitor do Kafka e o gravador do MaxCompute. Utilize Map Fields with Same Name, Map Fields in Same Line, Clear Mappings ou Manually Edit Mapping conforme necessário.

Campos padrão do Kafka

O Kafka expõe seis campos integrados que podem ser mapeados diretamente para colunas do MaxCompute.

Nome do campo

Descrição

__key__

A chave do registro do Kafka.

__value__

O valor do registro do Kafka.

__partition__

O número da partição onde o registro está localizado. Começa em 0.

__headers__

Os cabeçalhos do registro do Kafka.

__offset__

O offset do registro dentro de sua partição. Começa em 0.

__timestamp__

O timestamp do registro como um inteiro de 13 dígitos em milissegundos.

Análise de campos JSON

Para valores do Kafka formatados em JSON, adicione definições de campos personalizados usando . para acessar subcampos e [] para acessar elementos de array.

Dado este exemplo de valor de registro:

{
  "a": {
    "a1": "hello"
  },
  "b": "world",
  "c": [
    "xxxxxxx",
    "yyyyyyy"
  ],
  "d": [
    { "AA": "this", "BB": "is_data" },
    { "AA": "that", "BB": "is_also_data" }
  ],
  "a.b": "unreachable"
}

Definição do campo

Valor recuperado

Observações

a.a1

"hello"

Acesso a subcampo com .

b

"world"

Campo de nível superior

c[1]

"yyyyyyy"

Acesso a elemento de array com []

d[0].AA

"this"

Acesso combinado

a.b

(não recuperável)

Nomes de campos contendo . não podem ser analisados — o . é interpretado como separador de subcampo, causando ambiguidade.

Regras de mapeamento

  • A instância de sincronização não lê campos de origem não mapeados.

  • NULL é gravado em campos de destino não mapeados.

  • Um único campo de origem não pode ser mapeado para múltiplos campos de destino.

  • Um único campo de destino não pode receber mapeamento de múltiplos campos de origem.

Etapa 4: Configure parâmetros avançados

Clique em Advanced Configuration no lado direito da tarefa. Para este tutorial, defina Policy for Dirty Data Records como Ignore Dirty Data Records e mantenha todos os outros parâmetros com seus valores padrão. Para detalhes sobre os parâmetros, consulte Configuração de UI sem código.

Etapa 5: Testar a sincronização

  1. Clique em Run Configuration no lado direito da página de edição. Defina o Resource Group e os Script Parameters para a execução de teste e, em seguida, clique em Run na barra de ferramentas superior.

  2. Após a conclusão da execução de teste, verifique os dados na tabela de destino. No painel de navegação à esquerda, crie um arquivo com a extensão .sql e execute a seguinte consulta:

    Nota
    SELECT * FROM <MaxCompute_destination_table_name> WHERE pt=<specified_partition> LIMIT 20;

Etapa 6: Configure o agendamento e publicar

Clique em Scheduling no lado direito da tarefa para definir as configurações de agendamento da execução recorrente e, em seguida, clique em Publish para publicar a tarefa.

Os três parâmetros de agendamento usados neste tutorial — ${startTime}, ${endTime} e ${partition} — estão vinculados: os valores injetados em ${startTime} e ${endTime} controlam a janela de tempo do Kafka para cada instância, enquanto ${partition} determina qual partição do MaxCompute receberá os dados. Configure os três em conjunto com base no seu ciclo de agendamento.

Os exemplos a seguir cobrem os padrões de agendamento mais comuns.

Expressões de parâmetros de agendamento

Ciclo de agendamento

Expressão de startTime

Expressão de endTime

Expressão de partition

A cada 5 minutos

$[yyyymmddhh24mi-8/24/60]00

$[yyyymmddhh24mi-3/24/60]00

$[yyyymmddhh24mi-8/24/60]

A cada hora

$[yyyymmddhh24-1/24]0000

$[yyyymmddhh24]0000

$[yyyymmddhh24]

A cada 2 horas

$[yyyymmddhh24-2/24]0000

$[yyyymmddhh24]0000

$[yyyymmddhh24]

A cada 3 horas

$[yyyymmddhh24-3/24]0000

$[yyyymmddhh24]0000

$[yyyymmddhh24]

Diariamente

$[yyyymmdd-1]000000

$[yyyymmdd]000000

$[yyyymmdd-1]

Semanalmente

$[yyyymmdd-7]000000

$[yyyymmdd]000000

$[yyyymmdd-1]

Mensalmente

$[add_months(yyyymmdd,-1)]000000

$[yyyymmdd]000000

$[yyyymmdd-1]

Como as expressões funcionam (exemplo horário)

Para uma tarefa agendada para executar às 10:05 em 22/11/2022:

  • startTime resolve para 20221122090000 — lê registros a partir das 09:00 (inclusivo)

  • endTime resolve para 20221122100000 — para às 10:00 (exclusivo)

  • partition resolve para 2022112210 — grava nessa partição no MaxCompute

Exemplo de 5 minutos

Para uma tarefa agendada para executar às 10:00 em 22/11/2022:

  • startTime resolve para 20221122095200 — lê registros a partir das 09:52 (inclusivo)

  • endTime resolve para 20221122095700 — para às 09:57 (exclusivo)

  • partition resolve para 202211220952

O endTime é definido 3 minutos antes do horário de execução da instância. Essa margem garante que todos os registros da janela de tempo tenham sido gravados no Kafka antes que a instância comece a leitura.

Configurações do ciclo de agendamento

Ciclo de agendamento

Ciclo

Hora de início

Intervalo

Hora de término

Dia/data

A cada 5 minutos

Minuto

00:00

5 minutos

23:59

A cada hora

Hora

00:15

1 hora

23:59

Diariamente

Dia

00:15

Semanalmente

Semana

00:15

Segunda-feira

Mensalmente

Mês

00:15

1º de cada mês

Para agendamentos horários, diários, semanais e mensais, defina a hora de início logo após a meia-noite (por exemplo, 00:15 em vez de 00:00). Isso dá tempo ao Kafka para concluir a gravação de registros antes que a instância de sincronização inicie, reduzindo o risco de perda de dados.

Importante

Se registros com timestamps anteriores ou iguais a ${startTime} forem gravados no tópico do Kafka após uma instância recorrente já ter iniciado, esses registros poderão não ser lidos. Gravações atrasadas ou timestamps fora de ordem aumentam o risco de perda de dados.

Próximos passos

  • Revise a referência do Kafka Reader para obter a lista completa de opções de configuração.

  • Consulte Configuração de agendamento para opções avançadas de agendamento, como dependências e políticas de reexecução.