Use uma condição de filtro em um nó de sincronização em lote para sincronizar dados completos ou incrementais. Com essa condição, o Data Integration sincroniza apenas os dados que atendem aos critérios especificados. Combine parâmetros de agendamento com a condição de filtro para filtrar dados dinamicamente com base no tempo de execução do nó e viabilizar a sincronização incremental. Este tópico demonstra como configurar um nó de sincronização em lote para sincronização incremental.
Observações de uso
Algumas fontes de dados, como HBase e OTSStream, não oferecem suporte à sincronização incremental. Para verificar se uma fonte de dados específica suporta esse recurso, consulte a documentação do plug-in de leitura correspondente.
-
Os parâmetros obrigatórios para sincronização incremental variam conforme o plug-in de leitura. Para obter detalhes, consulte a documentação do plug-in específico e Fontes de dados e plug-ins suportados. Exemplo:
Plug-in de leitura
Parâmetro obrigatório
Sintaxe suportada
where
NotaNo modo assistente, este é o parâmetro de condição de filtro.
Sintaxe de banco de dados
NotaUse este parâmetro com parâmetros de agendamento para ler dados de um intervalo de tempo específico diariamente.
query
NotaNo modo assistente, este é o parâmetro Search Condition.
Semelhante à sintaxe de banco de dados
NotaUse este parâmetro com parâmetros de agendamento para ler dados de um intervalo de tempo específico diariamente.
Object
Especifique o caminho do objeto
NotaUse este parâmetro com parâmetros de agendamento para ler dados de um arquivo específico diariamente.
...
...
...
Configure a sincronização incremental
Em um nó de sincronização em lote do Data Integration, use parâmetros de agendamento para especificar caminhos de dados e intervalos nas tabelas de origem e destino. A configuração segue o mesmo padrão de outros tipos de nós.
Durante a execução, o sistema substitui todos os parâmetros de espaço reservado configurados no nó pelos valores reais representados pelas expressões dos parâmetros de agendamento e executa a sincronização dos dados.
Considere a sincronização de dados MySQL como exemplo:
Sem a configuração de Data Filtering, todos os dados são sincronizados para a tabela de destino por padrão.
Com a configuração de Data Filtering, apenas os dados que atendem à condição de filtro são sincronizados para a tabela de destino.
Um parâmetro de agendamento define o nome da partição da tabela MaxCompute de destino. $bizdate representa a data comercial. Quando uma tarefa agendada é executada, a expressão de partição configurada é substituída pela data comercial representada pelo parâmetro de agendamento. Para obter instruções detalhadas sobre expressões de parâmetros de agendamento, consulte Cenários de aplicação de parâmetros de agendamento no Data Integration. Em uma tarefa de sincronização em lote, configure o parâmetro bizdate em três locais para implementar a sincronização incremental: Na seção Data Filtering da origem, insira STR_TO_DATE('${bizdate}','%Y%m%d') <= gmt_modify_time AND gmt_modify_time < DATE_ADD(STR_TO_DATE('${bizdate}','%Y%m%d'), interval 1 day) para filtrar os dados modificados na data comercial. Na seção Partition Information do destino, insira pt=${bizdate} para gravar os dados na partição correspondente à data e defina a Cleanup Rule como Clean up existing data before writing (Insert Overwrite). Na seção Parameters das Schedule Settings, no lado direito, insira bizdate=$bizdate para que o sistema de agendamento substitua automaticamente ${bizdate} pela data comercial real durante a execução. Ao configurar a sincronização incremental de dados:
Sincronização incremental baseada em colunas temporais: Use parâmetros de agendamento para substituir dinamicamente dados temporais. Durante o agendamento da tarefa, esses parâmetros são substituídos automaticamente por valores específicos com base na data comercial. Para mais informações sobre parâmetros de agendamento, consulte Configure parâmetros de agendamento.
Sincronização incremental baseada em colunas não temporais: Use um nó de atribuição para converter a coluna para o tipo de dados alvo e passe-a para o Data Integration sincronizar. Para mais informações sobre nós de atribuição, consulte Crie um nó de atribuição.
Observações
Ao configurar uma tarefa de sincronização incremental, atente-se aos seguintes pontos:
Segurança da limpeza de dados existentes antes da gravação (Insert Overwrite): Quando múltiplas tarefas de sincronização gravam em partições diferentes da mesma tabela MaxCompute, a estratégia Insert Overwrite é segura. Essa estratégia limpa apenas os dados da partição especificada pela tarefa atual, sem afetar os dados em outras partições da tabela, evitando conflitos de dados ou exclusões acidentais.
Limitação de sobrescrita em lote por intervalo de partição: O DataWorks não suporta a especificação de um intervalo de horas (como
hh=00-23) na configuração de partição para sobrescrita em lote. Para sobrescrever dados de várias horas, configure tarefas separadas para cada hora. Atualmente, os parâmetros de partição suportam apenas um único valor específico ou o caractere curinga*.Sintaxe de caractere curinga: Quando a origem contém uma partição no nível de hora, mas o destino possui apenas uma partição no nível de dia, insira o caractere curinga
*no campo de partição de hora para corresponder a todos os dados horários. Insira*diretamente, sem aspas (como"*"). Caso contrário, ocorrerá um erro de sintaxe.
Sincronização incremental agendada de alta frequência baseada em timestamp
O DataWorks suporta sincronização incremental agendada baseada em timestamp combinando tarefas de sincronização em lote com agendamento periódico (por exemplo, a cada 5 minutos ou a cada hora). Essa abordagem é adequada para cenários de sincronização T+1 ou quase em tempo real do RDS MySQL para destinos como SelectDB e StarRocks. Ela implementa a sincronização incremental por meio de filtros SQL, sem exigir tarefas CDC em tempo real, o que evita os custos de execução contínua de tarefas. Pontos principais de configuração:
Na condição where da origem em Data Filtering, use uma coluna de timestamp como variável para filtragem. Exemplo:
gmt_modify_time >= '$[yyyymmddhhmiss-10/mi]' AND gmt_modify_time < '$[yyyymmddhhmiss]'.Configure parâmetros de agendamento periódico (como
$[yyyymmddhhmiss]) para calcular dinamicamente o intervalo de tempo e garantir que cada execução agendada sincronize apenas os dados incrementais dentro do intervalo especificado. Para obter detalhes sobre a configuração de parâmetros de agendamento, consulte Cenários de aplicação de parâmetros de agendamento no Data Integration.No mapeamento de colunas do destino, adicione manualmente um parâmetro constante mapeado para a coluna de partição para habilitar a gravação dinâmica de partições.
Configuração incremental para tarefas de sincronização em lote no nível de banco de dados
Além das tarefas de sincronização de tabela única, crie uma tarefa de sincronização em lote no nível de banco de dados para implementar a sincronização incremental periódica. Ao criar a tarefa, selecione a sincronização incremental e configure a condição incremental na tarefa de sincronização no nível de banco de dados. Isso permite uma sincronização eficiente de partições incrementais no nível de dia (por exemplo, filtrando pela coluna create_time).
Essa abordagem é ideal para cenários em que você deseja gerenciar centralmente a sincronização de várias tabelas, mas precisa de processamento incremental apenas para tabelas específicas. Para conhecer o processo completo de configuração de tarefas de sincronização em lote no nível de banco de dados, consulte Configure uma tarefa de sincronização em lote no nível de banco de dados.
Exemplos
Sincronizar dados históricos: Para sincronizar dados incrementais históricos para as partições de tempo correspondentes na tabela de destino, use o recurso de preenchimento de dados no Operation Center. Para mais informações sobre esse recurso, consulte Preencher dados. Na configuração do nó de sincronização de dados, selecione MySQL como fonte de dados e MaxCompute (ODPS) como destino, definindo o nome da tabela com um valor como
czd. Na condição de filtro de dados, use${bizdate}para controlar o intervalo incremental (por exemplo,STR_TO_DATE('${bizdate}','%Y%m%d') <= gmt_modify_time). Defina as informações de partição comods=${bizdate}e configure a regra de limpeza para Clean up existing data before writing (Insert Overwrite). Na seção de parâmetros das configurações de agendamento, definabizdate=$bizdate. Esse parâmetro de agendamento é substituído automaticamente pelo valor de data específico com base na data comercial durante o preenchimento. Ao executar o preenchimento, defina vários intervalos de datas comerciais (por exemplo, 01/05/2022 a 31/05/2022 e 01/04/2022 a 30/04/2022), selecione Immediately Run Backfill Instances Whose Scheduled Time Is Later Than the Current Time e escolha Ascending Order of Business Dates para a execução.Sincronização incremental de dados do ApsaraDB RDS para o MaxCompute
Perguntas frequentes
O que fazer se ocorrer um erro de partição não encontrada após o uso de parâmetros de agendamento no filtro de partição do MaxCompute Reader?
Causa: Os parâmetros de agendamento configurados não foram resolvidos corretamente para os valores reais de partição durante a execução, ou os valores resolvidos não correspondem às partições reais na tabela de origem.
Solução: Se o valor da partição for passado pelo parâmetro de saídas de um nó upstream, verifique a configuração de parâmetros no Data Studio e certifique-se de que as seguintes condições sejam atendidas:
O nome do parâmetro configurado deve ser consistente com o nome do parâmetro de entrada.
O valor do parâmetro passado deve corresponder exatamente à partição real no MaxCompute.
Uma tarefa de sincronização em lote do DataWorks executa sincronização completa ou incremental por padrão? Como configurar a sincronização incremental para uma tabela de origem sem colunas de partição?
Comportamento padrão
As tarefas de sincronização em lote do Data Integration do DataWorks realizam sincronização completa por padrão, o que significa que todos os dados são sincronizados a cada execução. A sincronização incremental só é ativada quando você configura uma condição de Data Filtering combinada com parâmetros de agendamento.
Tratamento de tabelas de origem sem colunas de partição
Se a tabela do banco de dados de origem (como uma tabela RDS) não possuir uma coluna de tempo ou de partição, não será possível filtrar dados incrementais diretamente usando uma condição where. Recomendamos adicionar uma coluna de tempo (como dt ou gmt_modify_time) à tabela de origem para servir como base para a filtragem incremental. Após preparar a coluna, configure a lógica de sincronização incremental seguindo a seção "Configure a sincronização incremental" deste tópico.
Conceitos principais
Definições de parâmetros
splitPk(chave de divisão): Especifica uma coluna de chave primária. O DataWorks divide os dados em vários blocos com base no intervalo de valores dessa coluna para permitir leituras simultâneas multithread.splitFactor(fator de divisão): Controla a granularidade da divisão. Um valor maior resulta em divisões mais refinadas e mais threads de leitura.
Impacto no desempenho e recomendações
Habilitar splitPk e splitFactor aumenta a carga no banco de dados de origem. Para reduzir essa carga, recomendamos:
Reduza a simultaneidade ou defina a simultaneidade da tarefa de sincronização em lote no nível de banco de dados como 1.
Garanta que a coluna de divisão (
splitPk) esteja indexada para melhorar a eficiência da leitura.