Ir para o conteúdo

Processar registros incrementais usando um high-watermark no Jitterbit Studio

Introdução

Um high-watermark é um valor armazenado que marca o ponto mais recente que uma operação de sincronização processou. A cada execução, a operação lê o watermark armazenado, usa-o para filtrar registros já processados e atualiza o watermark após a conclusão do lote. Apenas registros novos ou modificados são recuperados nas execuções subsequentes.

Este guia demonstra o padrão usando casos do Salesforce como fonte. O watermark é a data máxima de LastModifiedDate da execução anterior. O padrão se aplica a qualquer fonte que exponha um timestamp de modificação confiável e suporte consultas filtradas.

Duas funções de cache gerenciam o valor armazenado:

  • ReadCache: Recupera o watermark armazenado no início de cada execução.
  • WriteCache: Atualiza o watermark após a conclusão do processamento.

Para uma introdução ao ReadCache e WriteCache, incluindo opções de escopo de cache e limites de taxa, veja Detectar e deduplicar registros usando funções hash.

Padrão de design

Os passos de leitura e atualização do watermark cercam a lógica principal de consulta e processamento. A data armazenada avança após cada execução bem-sucedida, então a janela de consulta se desloca automaticamente para frente.

flowchart LR A["Script
Read watermark
from cache"] --> B["Script
Build filtered
SOQL query"] B --> C["Salesforce
Query activity
or SfLookupAll"] C --> D["Process records
(transformation or
child operations)"] D --> E["Script
Update watermark
in cache"]
Etapa Propósito
Ler watermark Ler a data armazenada; inicializar com um valor padrão se nada estiver armazenado.
Construir consulta filtrada Incorporar o watermark no filtro da consulta.
Consultar fonte Recuperar registros modificados na data do watermark ou depois dela.
Processar registros Executar a transformação ou cadeia de operações filhas.
Atualizar watermark Definir o watermark para a data máxima de modificação entre os registros recuperados e escrevê-lo no cache.

Parte 1: Ler o watermark

Adicione um passo de script como o primeiro passo da operação. O script lê a marca d'água armazenada e recorre a uma data padrão na primeira execução:

// Set the cache key and expiration
cacheKey = $project_name + "_LastModifiedDate";
cacheExpiry = 2592000; // 30 days in seconds

// Read the stored watermark
watermarkDate = ReadCache(cacheKey, cacheExpiry, "project");

// Fall back to the default date if no value is stored
if(length(trim(watermarkDate)) == 0,
    watermarkDate = $default_watermark_date;
);

Chave de cache: Use uma chave que seja única para este conjunto de dados dentro do projeto. Prefixar com o nome do projeto ou um identificador de conjunto de dados (por exemplo, $project_name + "_SF_Case_LM") evita colisões de chave quando várias marcas d'água são armazenadas no mesmo projeto.

Data padrão: default_watermark_date é uma variável de projeto definida em Variáveis de projeto para uma data passada suficientemente antiga para incluir todos os registros que você deseja na primeira execução. Use o formato ISO 8601 que o Salesforce espera no SOQL: por exemplo, 2000-01-01T00:00:00.000Z.

Expiração: 30 dias (2592000 segundos) mantém a marca d'água disponível entre execuções agendadas. Aumente esse valor para operações que são executadas com menos frequência.

Escopo: O escopo "project" torna o valor em cache acessível a todas as operações no projeto e o persiste entre as execuções. Use "env" se a marca d'água precisar ser compartilhada entre vários projetos no mesmo ambiente.

Parte 2: Consultar usando a marca d'água

Após o script ler a marca d'água, consulte a fonte usando watermarkDate como o limite do filtro.

Usando SfLookupAll em um script

SfLookupAll retorna um array bidimensional de registros correspondentes. Use-o quando o conjunto de resultados for passado para um script ou transformação subsequente para processamento:

soql = "SELECT Id, LastModifiedDate FROM Case"
    + " WHERE LastModifiedDate >= " + watermarkDate
    + " AND AccountId != null"
    + " ORDER BY LastModifiedDate ASC";

$caseIds = SfLookupAll("<TAG>endpoint:salesforce/Salesforce</TAG>", soql);

Para mais informações sobre como construir e executar consultas SOQL, veja Consultar registros do Salesforce usando SOQL.

Usando uma atividade de Consulta do Salesforce

Se você usar uma atividade de Consulta do Salesforce em vez de SfLookupAll, passe watermarkDate como uma variável do Jitterbit e faça referência a ela no campo de condição da atividade. Defina o operador de condição como maior ou igual a e insira [watermarkDate] como o valor.

Use um filtro maior ou igual e processamento idempotente

Um filtro estritamente maior (LastModifiedDate > a marca d'água) pode perder registros. Quando vários registros compartilham o mesmo timestamp de limite e apenas alguns foram incluídos no lote anterior, o filtro estrito exclui o restante na próxima execução, e eles nunca são processados. Use um filtro maior ou igual (>=) para que os registros de limite sejam re-incluídos. Isso recupera novamente os registros que definiram a marca d'água anterior, portanto, o processamento a montante deve ser idempotente (por exemplo, upsert por uma chave única ou deduplicar) para evitar a criação de duplicatas. Para uma abordagem de deduplicação, veja Detectar e deduplicar registros usando funções hash.

Parte 3: Atualizar a marca d'água

Após todos os registros no lote terem sido processados, defina a marca d'água como a máxima LastModifiedDate entre os registros realmente recuperados na Parte 2, e então escreva-a no cache. Adicione isso como o último passo do script na operação, ou na operação final da cadeia após todas as operações filhas serem concluídas:

// Derive the new watermark from the records retrieved in Part 2.
recordCount = Length($caseIds);
if(recordCount > 0,
    // Part 2 ordered results by LastModifiedDate ASC, so the last row holds the max.
    newWatermark = $caseIds[recordCount - 1]["LastModifiedDate"];
    WriteCache(cacheKey, newWatermark, cacheExpiry, "project");
);

Derive a marca d'água de registros processados, não do máximo da fonte

Não defina a marca d'água reconsultando a fonte para seu máximo atual (por exemplo, SELECT max(LastModifiedDate) FROM Case). Os registros podem ser modificados na fonte entre a consulta da Parte 2 e este passo de atualização. Esses registros não fazem parte do lote atual, mas uma consulta de máximo da fonte incluiria seus timestamps, avançando a marca d'água além de registros que nunca foram recuperados. Na próxima execução, o filtro estrito os ignora e eles são perdidos permanentemente. Sempre derive a marca d'água dos registros que esta execução realmente recuperou e processou.

Se você usar o caminho da atividade Query do Salesforce da Parte 2 em vez de SfLookupAll, capture a máxima LastModifiedDate à medida que os registros são processados (por exemplo, acumule o máximo em uma variável global na transformação de processamento), e então escreva esse valor no cache aqui.

O if guard impede que o WriteCache sobrescreva a marca d'água armazenada quando nenhum registro foi recuperado nesta execução.

cacheKey e cacheExpiry devem corresponder aos valores usados na Parte 1. Se o script de atualização for executado em uma operação diferente do script de leitura, atribua esses valores novamente ou armazene-os como variáveis de projeto para que ambos os scripts façam referência à mesma chave.

Verifique a integração

  1. Defina default_watermark_date para uma data de vários meses no passado. Implante e execute a operação. Adicione chamadas WriteToOperationLog para registrar o valor de watermarkDate e o número de registros retornados. Confirme que watermarkDate é igual ao padrão e que a consulta retornou registros.

  2. Execute a operação uma segunda vez sem modificar nenhum registro de origem. Como o filtro é maior ou igual, o registro (ou registros) cujo LastModifiedDate é igual à marca d'água armazenada é recuperado novamente. Confirme que apenas esses registros de limite são retornados (não o conjunto completo), o que significa que a marca d'água foi escrita corretamente após a primeira execução, e que o processamento idempotente não cria saída duplicada.

  3. Modifique um registro de origem e execute a operação novamente. Confirme que o registro modificado é retornado e processado. Registros que ainda estão no limite da marca d'água anterior também podem ser recuperados novamente; o processamento idempotente garante que eles não produzam duplicatas.

  4. Se a marca d'água não estiver persistindo entre as execuções, confirme que:

    • A chave de cache é idêntica em ambos os scripts de leitura e atualização.
    • Ambos os scripts usam o mesmo escopo ("project").
    • O valor escrito pelo WriteCache está no formato ISO 8601 que o Salesforce espera no SOQL.
  5. Se todos os registros forem recuperados em cada execução, a chamada WriteCache pode não estar sendo executada. Confirme que o script de atualização é executado após todas as operações filhas serem concluídas. Se você usar RunOperation para encadear operações, coloque o script de atualização na operação pai após a chamada RunOperation retornar, e não dentro de um loop de transformação.