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:
Fontes de dados do Kafka e do MaxCompute configuradas. Consulte Configuração de fonte de dados
Conectividade de rede entre o grupo de recursos e ambas as fontes de dados. Consulte Visão geral das soluções de conexão de rede
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
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 |
|
Read End Offset e End Time |
Defina Read End Offset como Specific Time e, em seguida, defina End Time como |
|
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 |
|
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 |
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 |
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 |
|
|
A chave do registro do Kafka. |
|
|
O valor do registro do Kafka. |
|
|
O número da partição onde o registro está localizado. Começa em 0. |
|
|
Os cabeçalhos do registro do Kafka. |
|
|
O offset do registro dentro de sua partição. Começa em 0. |
|
|
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 |
|
|
|
Acesso a subcampo com |
|
|
|
Campo de nível superior |
|
|
|
Acesso a elemento de array com |
|
|
|
Acesso combinado |
|
|
(não recuperável) |
Nomes de campos contendo |
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
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.
-
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
.sqle execute a seguinte consulta:NotaPara consultar dados dessa forma, anexe o projeto de destino do MaxCompute ao DataWorks como um recurso de computação.
Na página de edição do arquivo
.sql, clique em Run Configuration à direita, especifique o Computing Resource e o Resource Group e, em seguida, clique em Run.
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 |
|
|
|
|
A cada hora |
|
|
|
|
A cada 2 horas |
|
|
|
|
A cada 3 horas |
|
|
|
|
Diariamente |
|
|
|
|
Semanalmente |
|
|
|
|
Mensalmente |
|
|
|
Como as expressões funcionam (exemplo horário)
Para uma tarefa agendada para executar às 10:05 em 22/11/2022:
startTimeresolve para20221122090000— lê registros a partir das 09:00 (inclusivo)endTimeresolve para20221122100000— para às 10:00 (exclusivo)partitionresolve para2022112210— grava nessa partição no MaxCompute
Exemplo de 5 minutos
Para uma tarefa agendada para executar às 10:00 em 22/11/2022:
startTimeresolve para20221122095200— lê registros a partir das 09:52 (inclusivo)endTimeresolve para20221122095700— para às 09:57 (exclusivo)partitionresolve para202211220952
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.
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.