Pular para o conteúdo principal

Copie várias tabelas incrementalmente com uma tarefa For each

Quando você precisa copiar dados de muitas tabelas de origem para tabelas do Unity Catalog em uma programação, copiar todas as linhas em cada execução é lento e caro. Use uma *marca d'água* para rastrear a última linha processada para cada tabela e copiar apenas novas linhas em cada execução.

Este tutorial mostra como criar um Job orientado por metadados que:

  • Armazena a lista de tabelas de origem e o estado da marca d'água em uma tabela de controle Delta
  • Usa uma tarefa For each para processar cada tabela em paralelo
  • Copia apenas as linhas adicionadas desde a última execução bem-sucedida
  • Atualiza a marca d'água após cada cópia bem-sucedida.

Como funciona

O Job usa três tipos de tarefa conectados em sequência:

Tarefa

Tipo

O que faz

read_watermarks

SQL

Lê a tabela de controle de marca d'água e retorna uma linha por tabela de origem

copy_tables

Para cada

Itera sobre {{tasks.read_watermarks.output.rows}}, executando a tarefa aninhada uma vez por tabela de origem

copy_incremental (aninhado)

Notebook

Lê linhas adicionadas desde a última marca d'água, as grava na tabela de destino e avança a marca d'água

Tarefa

Tipo

O que faz

read_watermarks

SQL

Lê a tabela de controle de marca d'água e retorna uma linha por tabela de origem

copy_tables

Para cada

Itera sobre {{tasks.read_watermarks.output.rows}}, executando a tarefa aninhada uma vez por tabela de origem

copy_incremental (aninhado)

Notebook

Lê linhas adicionadas desde a última marca d'água, as grava na tabela de destino e avança a marca d'água

A saída da tarefa SQL — um array JSON de objetos de linha — flui para o campo Entradas da tarefa For each usando {{tasks.read_watermarks.output.rows}}. O Notebook aninhado recebe source_table, target_table, watermark_column e last_watermark para cada iteração.

Pré-requisitos

  • Um workspace Databricks com permissão para criar Job e Notebook
  • Permissão para criar esquemas e tabelas no Unity Catalog
  • Um SQL warehouse para executar tarefas de SQL

Este tutorial lê do dataset de exemplo samples.wanderbricks e grava em um esquema nomeado pelas variáveis catalog e schema no topo de cada bloco de código. Essas variáveis têm como default main.example_output. Para gravar em outro lugar, altere ambos os valores de forma consistente em cada bloco. O exemplo cria o esquema se ele não existir.

Etapa 1: Criar a tabela de controle de marca d'água

A tabela de controle de marca d'água é a fonte da verdade para quais tabelas processar e até onde cada tabela foi copiada. Cada linha representa uma tabela de origem.

Execute o seguinte SQL para criar a tabela de controle e registrar duas tabelas de origem. O exemplo registra samples.wanderbricks.users e samples.wanderbricks.properties, copiando cada um para uma tabela de destino em seu próprio esquema. As variáveis catalog e schema definem tanto o local da tabela de controle quanto os nomes da tabela de destino:

SQL
DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

CREATE SCHEMA IF NOT EXISTS IDENTIFIER(catalog || '.' || schema);

CREATE OR REPLACE TABLE IDENTIFIER(catalog || '.' || schema || '.watermarks') (
source_table STRING NOT NULL,
target_table STRING NOT NULL,
watermark_column STRING NOT NULL,
last_watermark TIMESTAMP NOT NULL
);

INSERT INTO IDENTIFIER(catalog || '.' || schema || '.watermarks') VALUES
('samples.wanderbricks.users', catalog || '.' || schema || '.users', 'created_at', '1970-01-01'),
('samples.wanderbricks.properties', catalog || '.' || schema || '.properties', 'created_at', '1970-01-01');

Ambas as tabelas de origem usam created_at como a coluna de watermark. Este é um Timestamp de tempo de inserção que apenas aumenta à medida que novas linhas chegam. Definir last_watermark como 1970-01-01 na primeira execução faz com que o notebook copie todas as linhas existentes. Isso atua como uma carga completa inicial. As execuções subsequentes copiam apenas as linhas adicionadas após a execução anterior.

nota

Recriar a tabela de controle redefine cada last_watermark para 1970-01-01, mas não limpa as tabelas de destino. Se você executar novamente este passo e, em seguida, executar o job novamente, o notebook tratará ambas as fontes como nunca copiadas e anexará cada linha histórica uma segunda vez. Para recomeçar de forma limpa, descarte as tabelas de destino ao recriar a tabela de controle.

O passo 2: Escreva o Notebook de cópia

O Notebook executa uma vez por iteração da tabela. Ele lê a marca d'água, filtra a origem, grava no destino e avança a marca d'água.

Crie um notebook em um caminho como /Workspace/Users/<username>/copy_incremental e adicione o código a seguir. Os padrões de widget permitem que você execute e teste o notebook diretamente. Quando executado dentro da tarefa For each, o job os substitui pelos valores de cada iteração, incluindo catalog e schema que localizam a tabela de controle.

O código lê apenas as linhas adicionadas desde o último watermark e as anexa à tabela de destino, criando-a caso ela não exista. Em seguida, ele calcula o high water mark das linhas que acabou de gravar e avança a tabela de controle para que a próxima execução comece a partir daí:

Python
from pyspark.sql.functions import max as spark_max

# Widget defaults let you run the notebook directly; the For each task overrides them per iteration
dbutils.widgets.text("catalog", "main", "Catalog")
dbutils.widgets.text("schema", "example_output", "Schema")
dbutils.widgets.text("source_table", "samples.wanderbricks.users", "Source table")
dbutils.widgets.text("target_table", "main.example_output.users", "Target table")
dbutils.widgets.text("watermark_column", "created_at", "Watermark column")
dbutils.widgets.text("last_watermark", "1970-01-01", "Last watermark")

catalog = dbutils.widgets.get("catalog")
schema = dbutils.widgets.get("schema")
source_table = dbutils.widgets.get("source_table")
target_table = dbutils.widgets.get("target_table")
watermark_column = dbutils.widgets.get("watermark_column")
last_watermark = dbutils.widgets.get("last_watermark")

# Read only rows newer than the last watermark. A strict > can skip rows that share the
# stored high-water timestamp; for insert-only sources with distinct timestamps this is safe.
new_rows = spark.table(source_table).filter(f"{watermark_column} > '{last_watermark}'")

row_count = new_rows.count()
print(f"Copying {row_count} new rows from {source_table}")

if row_count > 0:
# Append the new rows, creating the target table on the first run
new_rows.write.format("delta").mode("append").saveAsTable(target_table)

# Compute the high-water mark from the rows just written
new_watermark = new_rows.agg(spark_max(watermark_column)).collect()[0][0]

# Advance the control table so the next run starts from here
spark.sql(f"""
UPDATE {catalog}.{schema}.watermarks
SET last_watermark = CAST('{new_watermark}' AS TIMESTAMP)
WHERE source_table = '{source_table}'
""")

print(f"Watermark for {source_table} advanced to {new_watermark}")
else:
print(f"No new rows for {source_table}, watermark unchanged")
nota

Este notebook usa o modo append, que é adequado quando a origem contém apenas inserções, como samples.wanderbricks.users e samples.wanderbricks.properties fazem. Se a sua origem contiver atualizações, use watermark no Timestamp de atualização e use uma instrução MERGE em vez de write.mode("append") para fazer upsert de linhas na tabela de destino. Consulte Upsert em uma tabela Delta Lake usando merge para a sintaxe de merge.

Passo 3: Criar o Job

No seu workspace do Databricks, clique em fluxos de trabalho na barra lateral e depois clique em Criar Job . Atribua ao Job um nome como Incremental table copy.

Passo 4: Configure a tarefa de pesquisa de marca d'água

A tarefa SQL lê a tabela de controle e disponibiliza o resultado para a tarefa For each. Como a query declara as variáveis catalog e schema que localizam a tabela de controle, ela deve ser executada como um arquivo SQL de múltiplas instruções, em vez de no campo SQL inline da tarefa.

  1. Crie um arquivo SQL em seu workspace, como /Workspace/Users/<username>/read_watermarks.sql, com o seguinte conteúdo. Defina catalog e schema com os mesmos valores que você usou no passo 1:

    SQL
    DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
    DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

    SELECT source_table, target_table, watermark_column, last_watermark
    FROM IDENTIFIER(catalog || '.' || schema || '.watermarks');
  2. No job, clique em Adicionar tarefa .

  3. Definir Nome da tarefa como read_watermarks.

  4. Defina Type como SQL e, em seguida, defina SQL tarefa como File .

  5. Defina Caminho para o arquivo SQL que você criou.

  6. Defina **SQL warehouse** para um warehouse em seu workspace.

  7. Clique em Criar tarefa .

Quando esta tarefa é executada, o Databricks captura o resultado como uma matriz JSON em tasks.read_watermarks.output.rows. Após uma carga completa inicial, cada last_watermark reflete a linha mais recente copiada dessa origem:

JSON
[
{
"source_table": "samples.wanderbricks.users",
"target_table": "main.example_output.users",
"watermark_column": "created_at",
"last_watermark": "2025-07-30T23:05:18.000Z"
},
{
"source_table": "samples.wanderbricks.properties",
"target_table": "main.example_output.properties",
"watermark_column": "created_at",
"last_watermark": "2025-07-30T00:00:00.000Z"
}
]

Passo 5: Configurar a tarefa For each

A tarefa For each lê a saída SQL e inicia uma execução de tarefa aninhada por tabela de origem.

  1. Clique em Adicionar tarefa e defina Depende de como read_watermarks.

  2. Definir Nome da tarefa como copy_tables.

  3. Defina o **Tipo** como **Para cada**.

  4. No campo **Entradas**, insira:


    {{tasks.read_watermarks.output.rows}}
  5. Defina **Concorrência** como 2 para copiar duas tabelas por vez. Aumente este valor se seu warehouse puder suportar maior paralelismo.

  6. Clique em **Adicionar uma tarefa para fazer loop** para configurar a tarefa aninhada.

  7. Definir Nome da tarefa como copy_incremental.

  8. Set Type to Notebook .

  9. Defina **Caminho** como o caminho do Notebook que você criou no Passo 2.

  10. Clique em Parâmetros , então clique em Adicionar para adicionar cada um dos seguintes parâmetros:

Chave

Valor

catalog

main

schema

example_output

source_table

{{input.source_table}}

target_table

{{input.target_table}}

watermark_column

{{input.watermark_column}}

last_watermark

{{input.last_watermark}}

Chave

Valor

catalog

main

schema

example_output

source_table

{{input.source_table}}

target_table

{{input.target_table}}

watermark_column

{{input.watermark_column}}

last_watermark

{{input.last_watermark}}

Defina catalog e schema com os mesmos valores que você usou no passo 1 para que o notebook avance a tabela de controle que a tarefa SQL lê. Cada referência {{input.<key>}} é resolvida para o campo correspondente da linha da iteração atual. 11. Clique em Criar tarefa .

Passo 6: Execute o Job e verifique

  1. Clique em Executar agora para acionar o Job.
  2. Na página de execução do Job, clique no nó copy_tables para expandir a tarefa For each.
  3. A página de execução mostra uma tabela de iterações — uma linha por tabela de origem — cada uma exibindo seu status, horário de início e duração.
  4. Clique em qualquer iteração para exibir a saída do Notebook e confirmar a contagem de linhas e a atualização da marca d'água.

Para confirmar que a marca d'água avançou, execute a seguinte query após a conclusão do job. Defina catalog e schema com os mesmos valores que você usou nas etapas anteriores:

SQL
DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

SELECT source_table, last_watermark
FROM IDENTIFIER(catalog || '.' || schema || '.watermarks');

Cada valor last_watermark agora deve refletir o carimbo de data/hora da linha copiada mais recentemente. Se um valor ainda for 1970-01-01, a tabela de origem não continha linhas correspondentes ao filtro, ou a tarefa de cópia encontrou um erro — verifique a saída da execução da tarefa para detalhes.

Estender o padrão

Cada snippet declara as mesmas variáveis catalog e schema usadas nos passos anteriores. Defina-os com os valores que localizam sua tabela de controle.

Adicionar uma nova tabela de origem : insira uma linha na tabela de controle. A próxima execução do job a captura automaticamente, começando com uma carga completa a partir de 1970-01-01. Este snippet pressupõe que você já adicionou a coluna active de Pause a table abaixo, portanto, ele define active como true. As duas extensões podem ser aplicadas em qualquer ordem; se você ainda não adicionou a coluna, descarte o último valor e sua coluna:

SQL
DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

INSERT INTO IDENTIFIER(catalog || '.' || schema || '.watermarks')
(source_table, target_table, watermark_column, last_watermark, active)
VALUES
('samples.wanderbricks.hosts', catalog || '.' || schema || '.hosts', 'joined_at', '1970-01-01', TRUE);

Pausa uma tabela : adicione uma coluna active, preencha-a com true para as linhas existentes e, em seguida, filtre-a na tarefa de arquivo SQL. O Delta requer a adição da coluna e a definição de seu valor em instruções separadas:

SQL
DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

ALTER TABLE IDENTIFIER(catalog || '.' || schema || '.watermarks') ADD COLUMN active BOOLEAN;

UPDATE IDENTIFIER(catalog || '.' || schema || '.watermarks') SET active = TRUE;

Em seguida, adicione WHERE active = TRUE ao SELECT em seu arquivo read_watermarks.sql para que o job ignore as tabelas em pausa.

Retroalimentar uma tabela : Reset sua marca d'água para copiar novamente de um ponto específico:

SQL
DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

UPDATE IDENTIFIER(catalog || '.' || schema || '.watermarks')
SET last_watermark = '2025-01-01'
WHERE source_table = 'samples.wanderbricks.users';

Recursos adicionais