As tarefas de sincronização em tempo real de tabela única no DataWorks Data Integration permitem replicar dados com baixa latência e alto throughput entre origens de dados. O mecanismo de computação em tempo real captura alterações nos dados de origem (inserções, exclusões e atualizações) e as aplica ao destino. Este tópico usa a sincronização em tempo real do Kafka para o MaxCompute como exemplo.
Preparações
-
Preparar origens de dados
Crie as origens de dados de origem e de destino. Para mais informações, consulte Gerenciamento de origens de dados.
Verifique se as origens de dados suportam sincronização em tempo real. Para mais informações, consulte Origens de dados e soluções de sincronização suportadas.
Algumas origens de dados, como Hologres e Oracle, exigem a ativação de logs. O método varia conforme a origem de dados. Para mais informações, consulte Lista de origens de dados.
Grupo de recursos: Adquira e configure um grupo de recursos Serverless.
Conectividade de rede: Estabeleça uma conexão de rede entre o grupo de recursos e as origens de dados.
Etapa 1: Criar uma tarefa de sincronização
Faça login no console do DataWorks. Na região de destino, clique em no painel de navegação à esquerda. Selecione um workspace na lista suspensa e clique em Go to Data Integration.
-
No painel de navegação à esquerda, clique em Synchronization Task. No topo da página, clique em Create Synchronization Task e configure a tarefa. Este exemplo usa a sincronização em tempo real do Kafka para o MaxCompute:
Source Type:
Kafka.Destination Type:
MaxCompute.Specific Type:
Single-table real-time.-
Synchronization Mode:
Schema Migration: Cria automaticamente objetos de banco de dados no destino, como tabelas, campos e tipos de dados, correspondentes à origem. Esta etapa não transfere dados.
-
Incremental Sync (opcional): Após a conclusão da sincronização completa, esta etapa captura continuamente as alterações de dados (inserções, atualizações e exclusões) da origem e as aplica ao destino.
Se a origem for o Hologres, a sincronização completa também é suportada. Esse processo sincroniza primeiro todos os dados existentes para a tabela de destino. Em seguida, a sincronização incremental de dados inicia automaticamente.
Para mais informações sobre origens de dados e soluções de sincronização suportadas, consulte Origens de dados e soluções de sincronização suportadas.
Etapa 2: Configurar origens de dados e recursos de execução
Em Source Information, selecione sua origem de dados
Kafka. Em Destination, selecione sua origem de dadosMaxCompute.Na seção Running Resources, selecione o Resource Group para a tarefa de sincronização e atribua Resource Group CUs à tarefa. Defina CUs separadamente para sincronização completa e incremental para controlar os recursos com precisão e evitar desperdício. Se a tarefa de sincronização falhar com um erro de falta de memória (OOM), aumente o valor de CU do grupo de recursos.
Verifique se ambas as origens de dados (origem e destino) passaram na Connectivity Check.
Etapa 3: Configurar a solução de sincronização
1. Configurar a origem
-
Na aba Configuration, selecione o tópico do Kafka a ser sincronizado.
Use os valores padrão para outras configurações ou modifique-os conforme necessário. Para mais informações, consulte a documentação oficial do Kafka.
-
No canto superior direito, clique em Data Sampling.
Na caixa de diálogo exibida, defina o Start time e os Sampled Data Records e clique em Start Collection. Essa ação amostra dados do tópico Kafka especificado. Visualize os dados amostrados, que servem como entrada para as configurações de visualização e preview dos nós subsequentes de processamento de dados.
-
Na aba Configure Output Field, selecione os campos a serem sincronizados.
Por padrão, o Kafka fornece seis campos.
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 do Kafka está localizado. O número da partição é um inteiro que começa em 0.
__headers__
Os cabeçalhos do registro do Kafka.
__offset__
O offset do registro do Kafka em sua partição. O offset é um inteiro que começa em 0.
__timestamp__
O timestamp UNIX de 13 dígitos em milissegundos para o registro do Kafka.
Também é possível transformar campos adicionalmente nos nós subsequentes de processamento de dados.
2. Processamento de dados
Ative a opção Data Processing. Cinco métodos de processamento de dados estão disponíveis: data masking, string replace, data filtering, JSON parsing e edit and assign fields. Organize esses métodos em qualquer ordem. Durante a execução, eles são processados sequencialmente.
Após configurar um nó de processamento de dados, clique em Preview Data Output no canto superior direito:
A tabela abaixo dos dados de entrada mostra os resultados da etapa anterior de Data Sampling. Clique em Re-obtain Output of Ancestor Node para atualizar os resultados.
Se não houver saída upstream, use Manually Construct Data para simular a saída anterior.
Clique em Preview para visualizar a saída processada pelo componente de processamento de dados.
A tabela Input Data exibe campos de mensagens do Kafka como key, value, partition, offset e timestamp. A seção Preview Result mostra os campos analisados após o processamento de dados, como pares chave-valor resultantes da análise de JSON. Ela também indica o número de registros de dados sujos e inclui um aviso: "O preview mostrado aqui é apenas para referência. Os resultados da execução real da tarefa são os que valem."
O preview de saída de dados e os recursos de processamento de dados dependem da Data Sampling da origem Kafka. Execute a amostragem de dados na origem Kafka antes de processar os dados.
3. Configurar o destino
Na seção Destination, selecione um grupo de recursos Tunnel. A seleção padrão é "Public Transport Resources", que corresponde à cota gratuita do MaxCompute.
-
Especifique se deseja gravar em uma nova tabela ou usar uma tabela existente.
Ao criar uma nova tabela, selecione Create na lista suspensa. Por padrão, o sistema cria automaticamente uma tabela com a mesma estrutura da origem. Altere manualmente o nome e o schema da tabela de destino, se necessário.
Caso opte por usar uma tabela existente, selecione a tabela de destino na lista suspensa.
-
(Opcional) Edite o schema da tabela.
Clique no ícone de edição ao lado do nome da tabela para modificar seu schema. Clique em Re-generate Table Schema Based on Output Column of Ancestor Node para gerar automaticamente o schema com base nas colunas de saída do nó upstream. Em seguida, selecione uma coluna no schema gerado automaticamente para ser a chave primária.
4. Configurar mapeamento de campos
Após selecionar a origem e o destino, especifique o mapeamento entre as colunas de origem e de destino. A tarefa grava dados dos campos de origem nos campos de destino correspondentes com base nesse mapeamento.
O sistema mapeia automaticamente as colunas upstream para as colunas de destino com base no princípio de The same name mapping. Ajuste os mapeamentos conforme necessário. Uma única coluna upstream pode ser mapeada para várias colunas de destino, mas múltiplas colunas upstream não podem ser mapeadas para uma única coluna de destino. Se uma coluna upstream não for mapeada, seus dados não serão gravados na tabela de destino.
-
Para campos do Kafka, configure uma análise de JSON personalizada. Utilize o componente de processamento de dados para recuperar o conteúdo do campo value e obter uma configuração de campos mais granular.
Na seção de mapeamento de campos, clique em Map by Same Name para criar uma correspondência um a um entre os campos de entrada upstream e os campos de saída do MaxCompute. O mapeamento inclui campos de metadados do Kafka (
_key_,_value_,_partition_,_offset_,_timestamp_e_headers_) e campos de negócios (id,nameeage). Campos do tipo LONG são mapeados para o tipo BIGINT no MaxCompute, e os tipos STRING permanecem inalterados. -
(Opcional) Configure partições.
O Particionamento automático baseado em tempo cria partições com base no tempo de negócios (neste caso, o campo _timestamp). A partição de primeiro nível é por ano, a de segundo nível por mês, e assim por diante.
O Particionamento dinâmico por conteúdo de campo mapeia um campo da tabela de origem para um campo de partição na tabela de destino do MaxCompute. Isso garante que linhas contendo dados específicos no campo de origem sejam gravadas na partição correspondente na tabela do MaxCompute.
Etapa 4: Configuração avançada
A tarefa de sincronização oferece parâmetros avançados para configuração detalhada. Valores padrão são fornecidos e, na maioria dos casos, não precisam ser alterados. Para modificá-los:
-
No canto superior direito da página, clique em Advanced Settings para acessar a página de configuração de Advanced Parameters.
NotaNo Data Development, a configuração avançada encontra-se em uma aba no lado direito da página de configuração da tarefa.
Defina parâmetros separadamente para as extremidades de leitura e gravação da tarefa de sincronização. Defina Automatically set runtime configuration como false para personalizar a Runtime Configuration.
Modifique os valores dos parâmetros com base nas dicas de ferramenta. A descrição de cada parâmetro é exibida ao lado de seu nome. Para sugestões de configuração de alguns parâmetros, consulte Parâmetros avançados para sincronização em tempo real.
Modifique esses parâmetros somente após compreender totalmente sua finalidade e possíveis consequências. Configurações incorretas podem causar erros inesperados ou problemas de qualidade de dados.
Etapa 5: Execução de teste
Após concluir todas as configurações da tarefa, clique em Perform Simulated Running no canto inferior esquerdo para depurar a tarefa. Isso simula como a tarefa processa uma pequena quantidade de dados de amostra e permite visualizar os resultados na tabela de destino. Se houver erros de configuração, exceções ou dados sujos, o sistema fornece feedback em tempo real para ajudar a verificar a configuração da tarefa.
Na caixa de diálogo exibida, defina os parâmetros de amostragem: Start time e Sampled Data Records.
Clique em Start Collection para recuperar os dados de amostra.
Clique em Preview Result para simular a execução da tarefa e visualizar a saída.
A saída da execução de teste serve apenas para preview. Ela não é gravada na origem de dados de destino e não afeta os dados de produção.
Etapa 6: Publicar e executar
Após concluir todas as configurações, clique em Save na parte inferior da página para salvar a configuração da tarefa.
As tarefas do Data Integration devem ser publicadas no ambiente de produção para serem executadas. Qualquer tarefa nova ou editada precisa ser submetida via Deploy para entrar em vigor. Durante a publicação, se você selecionar Start immediately after deployment, a tarefa iniciará automaticamente. Caso contrário, após a publicação da tarefa, acesse a página e inicie a tarefa manualmente na coluna Actions.
Em Tasks, clique no Name/ID da tarefa para visualizar seu processo detalhado de execução.
Etapa 7: Configurar regras de alerta
Depois que a tarefa for publicada e estiver em execução, configure regras de alerta para receber notificações imediatas sobre exceções. Isso ajuda a manter a estabilidade do ambiente de produção e a atualidade dos dados. Na lista de tarefas do Data Integration, localize a tarefa de destino e, na coluna Actions, clique em .
1. Adicionar um alerta
Selecione Use custom rules e insira um Alert Name e uma Description. Os métodos de notificação suportados incluem Email, SMS, Phone, DingTalk, webhook e Feishu. Configure esses métodos independentemente para os níveis WARNING e CRITICAL. Especifique os destinatários por meio de um on-call schedule ou adicione-os manualmente.
(1) Clique em Create Rule para configurar uma regra de alerta.
Defina o Alert Reason para monitorar métricas como Business delay, failover, Task status, DDL Notification e Task Resource Utilization. Em seguida, defina níveis de alerta CRITICAL ou WARNING com base em limiares especificados.
Após definir o método de alerta, use Configure Advanced Parameters para controlar o intervalo de envio de notificações de alerta e evitar sobrecarga de mensagens.
Se selecionar Business delay, Task status ou Task Resource Utilization como causa do alerta, ative também as notificações de recuperação para informar aos destinatários quando a tarefa retornar ao estado normal.
(2) Gerenciar regras de alerta.
Para regras de alerta existentes, use a opção alternar para ativá-las ou desativá-las. Também é possível enviar alertas para diferentes pessoas com base no nível do alerta.
2. Visualizar alertas
Na lista de tarefas, expanda para acessar a página de eventos de alerta e visualizar alertas históricos.
Mais operações
Após o início da tarefa, clique no nome da tarefa para visualizar seus detalhes de execução e realizar O&M e ajuste de tarefas.
FAQ
Para perguntas frequentes sobre tarefas de sincronização em tempo real, consulte FAQ sobre sincronização em tempo real.
Referência
Sincronização em tempo real de tabela única do Kafka para o ApsaraDB for OceanBase
Ingestão em tempo real de tabela única do LogHub (SLS) para o Data Lake Formation
Sincronização em tempo real de tabela única do Hologres para o Doris
Sincronização em tempo real de tabela única do Hologres para o Hologres
Sincronização em tempo real de tabela única do Kafka para o Hologres
Sincronização em tempo real de tabela única do LogHub (SLS) para o Hologres
Sincronização em tempo real de tabela única do Kafka para o Hologres
Sincronização em tempo real de tabela única do Hologres para o Kafka
Sincronização em tempo real de tabela única do LogHub (SLS) para o MaxCompute
Sincronização em tempo real de tabela única do Kafka para um data lake OSS
Sincronização em tempo real de tabela única do Kafka para o StarRocks
Sincronização em tempo real de tabela única do Oracle para o Tablestore