Controlo de concorrência para tabelas Delta

Quando vários notebooks do Fabric, pipelines ou trabalhos do Spark escrevem na mesma tabela Delta ao mesmo tempo, o Delta Lake utiliza controlo de concorrência otimista (OCC) para manter a tabela consistente. Cada transação lê uma captura instantânea, escreve novos ficheiros e depois verifica se não ocorreu nenhum commit em conflito entretanto. Se for detetado um conflito, a transação falha com uma exceção em vez de corromper os dados.

Este artigo aborda padrões práticos para gerir escritas simultâneas no Fabric. Para uma especificação completa do protocolo OCC, veja controlo de concorrência Delta Lake (documentação open-source).

Níveis de isolamento

Todas as tabelas Delta utilizam o nível de isolamento Serializable. Serializable é o nível mais rigoroso e o único suportado. Garante que o resultado das transações concorrentes é idêntico a uma ordem de execução sequencial.

O Delta Lake também utiliza um nível interno de SnapshotIsolation para operações que não alteram dados a nível lógico (como OPTIMIZE). O SnapshotIsolation ignora a verificação de anexação concorrente, permitindo que a compactação prossiga sem conflito com inserções concorrentes. Não se configura o SnapshotIsolation diretamente — o Delta Lake aplica-o automaticamente quando apropriado.

Com isolamento Serializable, uma anexação cega concorrente (INSERT INTO) pode entrar em conflito com um MERGE ou UPDATE que lê a mesma partição.

Que operações entram em conflito

Nem todas as operações de escrita concorrentes entram em conflito. O fator chave é se duas operações tocam nos mesmos ficheiros subjacentes.

Par concorrente Conflito? Porquê
Duas operações de INSERT (acrescentar) No Cada um adiciona novos ficheiros sem ler os já existentes (anexação cega).
INSERT + OPTIMIZE No OPTIMIZE faz commits em SnapshotIsolation porque não altera dados lógicos, por isso ignora completamente a verificação de anexação concorrente. As anexações adicionam novos ficheiros que não se sobrepõem aos ficheiros que estão a ser compactados.
Duas UPDATE, DELETE ou MERGE operações Sim, se lerem ou modificarem ficheiros sobrepostos Cada um reescreve os ficheiros, por isso o instantâneo do segundo escritor está desatualizado.
OPTIMIZE + UPDATE/DELETE/MERGE Sim, se tocarem nos mesmos ficheiros OPTIMIZE remove e readiciona ficheiros (com dataChange=false). Se uma operação de modificação de dados também lesse esses mesmos ficheiros, é gerada uma ConcurrentDeleteReadException.
Duas OPTIMIZE corridas Sim, se selecionarem os mesmos ficheiros Ambos tentam remover e reescrever o mesmo conjunto de ficheiros, provocando um ConcurrentDeleteDeleteException.
INSERT + MERGE/UPDATE/DELETE Sim, se a operação de modificação de dados ler a mesma partição Em Serializable, um acrescento cego pode entrar em conflito com modificações de dados simultâneas se a operação tiver lido uma partição na qual o acrescento gravou dados.

Sugestão

Os pipelines só de acréscimo (INSERT INTO, df.write.mode("append")) são a forma mais simples de evitar totalmente os conflitos. Se a tua carga de trabalho puder ser adicionada primeiro e reconciliada depois, eliminas a contenção entre escrita e escrita.

Isolar escritores com particionamento

A forma mais comum de executar DML concorrente contra a mesma tabela sem conflitos é particionar a tabela pela coluna que separa os seus escritores, e depois incluir essa coluna em todas as condições de operação. Quando cada escritor tem como alvo uma partição diferente, as operações tocam conjuntos de ficheiros desjuntos e não entram em conflito.

Um cenário típico: vários pipelines processam, cada um, dados para diferentes unidades de negócio ou tenants. Particione por essa dimensão e afixe o MERGE de cada canalização à sua partição.

-- Each pipeline targets its own partition, so concurrent runs don't conflict
MERGE INTO events AS target
USING staged AS source
ON target.event_id = source.event_id
    AND target.business_unit = 'EMEA'
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *

Importante

A coluna de partição deve aparecer na própria condição de fusão — não apenas nos dados de origem. Sem ela, a Delta Lake não consegue determinar na altura da validação que as duas operações tocaram conjuntos de ficheiros desconexos, e o verificador de conflitos trata a operação como uma leitura de tabela completa.

Para mais detalhes sobre estratégias de particionamento, veja Particionamento para tabelas Delta.

Retentativa de confirmação integrada

O Delta Lake repete automaticamente a tentativa de efetuar um commit quando deteta que outra transação efetuou primeiro o commit. Em cada nova tentativa, lê o commit vencedor, executa o verificador de conflito e — se não existir conflito lógico — tenta novamente o commit na próxima versão disponível. Este processo repete-se de forma transparente, sem qualquer ação do seu código.

Um conflito lógico (por exemplo, duas operações a reescrever o mesmo ficheiro) não pode ser resolvido automaticamente. A nova tentativa gera uma das exceções listadas em Exceções comuns de conflito. No entanto, muitas colisões transitórias de versões, como dois acrescentos sem verificação prévia a disputar a mesma posição de versão, são resolvidas automaticamente e nem sequer chegam à sua aplicação.

Exceções de conflito comum

Quando é detetado um conflito, Delta Lake levanta uma exceção específica. Compreender qual a exceção que se vê ajuda a identificar a causa raiz.

Exception O que aconteceu
ConcurrentAppendException Outro escritor adicionou ficheiros numa partição (ou conjunto de ficheiros) que a sua operação estava a ler. Comum quando um MERGE é executado sobre uma partição que também está a receber operações de inserção por outro pipeline. Sob isolamento Serializable, até os acrescentos cegos (operações INSERT simples) podem desencadear esta exceção.
ConcurrentDeleteReadException Outro escritor apagou ou reescreveu um ficheiro que a sua operação leu. Isto é típico quando OPTIMIZE compacta ficheiros que também estavam a ser lidos por um UPDATE ou MERGE em simultâneo, ou quando duas operações de modificação de dados se sobrepõem sobre as mesmas linhas.
ConcurrentDeleteDeleteException Ambas as operações tentaram apagar ou reescrever o mesmo ficheiro. Muitas vezes causado por execuções sobrepostas OPTIMIZE ou por dois pipelines a reescreverem a mesma partição simultaneamente.
ConcurrentWriteException É gerado um conflito genérico quando outra transação foi confirmada na mesma versão da tabela antes de a resolução de conflitos poder ser executada — por exemplo, durante uma atualização do sistema de ficheiros para confirmações geridas.
MetadataChangedException O esquema ou as propriedades da tabela mudaram durante a transação, por exemplo, devido a uma escrita concorrente ALTER TABLE ou a uma operação de escrita de evolução do esquema.
ConcurrentTransactionException Duas consultas de Streaming Estruturado com a mesma localização do ponto de verificação gravaram na tabela ao mesmo tempo. Desduplica os teus trabalhos de streaming ou usa caminhos de checkpoint distintos.
ProtocolChangedException Uma transação simultânea atualizou ou rebaixou o protocolo da tabela enquanto a transação atual também tentava alterar o protocolo. Também pode ocorrer quando uma característica da tabela é eliminada simultaneamente.

Estratégias comuns para evitar conflitos de escrita

Ativar a compactação automática

A autocompactação corre de forma síncrona como parte das operações de escrita. A compactação síncrona impede que tarefas de compactação agendadas separadamente coincidam com operações de modificação de dados e, consequentemente, causem exceções de escrita simultânea.

Agendar manutenção fora das janelas de escrita

OPTIMIZE e VACUUM podem entrar em conflito com operações simultâneas de modificação de dados. No Fabric, agende trabalhos de notebook ou atividades do pipeline para compactação de tabelas e VACUUM durante janelas de baixa atividade — por exemplo, após a conclusão da ingestão noturna, em vez de durante a mesma.

Use padrões de anexação e fusão

Para a ingestão com elevada concorrência, grave os dados brutos com operações de escrita só de acréscimo numa tabela de preparação (sem possibilidade de conflitos) e, em seguida, execute uma única tarefa MERGE para reconciliar os dados na tabela de destino. O padrão serializa a operação suscetível a conflitos, mantendo a ingestão totalmente paralela.

Adicionar mecanismo de repetição para conflitos lógicos

A nova tentativa de commit integrada trata automaticamente conflitos transitórios de versão, mas os conflitos lógicos — quando duas operações se sobrepõem efetivamente — geram uma exceção. Como o Delta Lake nunca produz escritas parciais, uma transação falhada é segura para tentar novamente ao nível da aplicação. Para pipelines em que são esperados conflitos lógicos ocasionais, envolva a operação de escrita numa lógica de nova tentativa:

from delta.exceptions import ConcurrentAppendException
import time

# Retry with backoff on transient concurrent write conflicts
max_retries = 3
for attempt in range(max_retries):
    try:
        spark.sql("MERGE INTO target USING source ON ...")
        break
    except ConcurrentAppendException:
        if attempt < max_retries - 1:
            time.sleep(2 ** attempt)
        else:
            raise

Escolha a estratégia de esquema adequada

O agrupamento de líquidos e o particionamento resolvem diferentes problemas. O agrupamento de líquidos otimiza o layout dos ficheiros para o desempenho de leitura. A partição cria limites físicos que evitam conflitos concorrentes entre escritores. Se a sua carga de trabalho precisar de ambos, particione pela coluna de isolamento de escrita e utilize Z-Order dentro de cada partição para melhorar o desempenho de leitura.