Data Drains

Os Data Drains permitem que owners e admins de organizações em planos Enterprise exportem continuamente dados do Studio para um destino sob seu controle — um bucket S3, um bucket do Google Cloud Storage, um container do Azure Blob, uma tabela do BigQuery, uma tabela do Snowflake, a ingestão de logs do Datadog ou um webhook HTTPS, todos de propriedade do cliente. Um drain é executado em um agendamento, pega apenas as linhas novas desde a última execução bem-sucedida e as grava no destino. Ver a configuração do drain e o histórico de execuções também é restrito a owners e admins, já que os destinos expõem nomes internos de bucket, identificadores de tabela e URLs de webhook.

Os drains são independentes da Retenção de Dados, mas foram projetados para funcionar em conjunto — veja Combinando com a Retenção de Dados abaixo.


Configuração

Vá para Settings → Enterprise → Data Drains no seu workspace e clique em New drain.

Página de configurações de Data Drains mostrando dois drains configurados — um exportando logs de workflow para o Amazon S3 diariamente, outro exportando conversas do Chat para um webhook HTTPS a cada hora

Diálogo de novo data drain com campos para nome, fonte, cadência, destino e credenciais do S3

Cada drain tem quatro partes:

  1. Uma source — a categoria de dados a exportar
  2. Um destination — para onde os dados vão
  3. Um schedule — com que frequência é executado
  4. Um name — único dentro da sua organização

Fontes

Um drain exporta exatamente uma fonte. Para exportar várias fontes, crie vários drains.

FonteDescrição
Workflow logsRegistros de execução de workflow (uma linha por execução, somente depois que a execução atinge um estado terminal).
Job logsRegistros de jobs em segundo plano (APIs com deploy, agendamentos, webhooks). Somente linhas em estado terminal são exportadas.
Audit logsEventos de auditoria no escopo da organização e do workspace — logins, mudanças de permissões, criação/exclusão de recursos, mudanças na configuração de drains.
ChatsHistórico do Chat.
Chat runsRegistros de execuções do Chat (somente estado terminal).

Cada linha é entregue como uma única linha de NDJSON. O formato de cada linha faz parte do schema público e é estável entre versões; toda linha carrega um campo id que os consumidores downstream podem usar para deduplicar.

Os drains exportam cada linha exatamente uma vez, com base no cursor de criação. Campos mutáveis em Chats (mensagens, título, lastSeenAt) são um retrato de um momento e não serão emitidos novamente se o chat for atualizado depois. Trate a exportação como append-only e reconstrua o estado atual a partir do seu próprio sistema de registro, se precisar dele.


Destinos

Amazon S3 (ou qualquer store compatível com S3)

Grava um objeto NDJSON por chunk entregue no seu bucket.

  • Bucket — o nome do bucket. Precisa já existir; o Studio não cria buckets.
  • Region — região da AWS (por exemplo, us-east-1).
  • Prefix (opcional) — caminho da pasta dentro do bucket. A barra final é opcional.
  • Access key ID / Secret access key — credenciais IAM com s3:PutObject no bucket. O botão "Test connection" faz uma gravação real de teste para verificar e depois a exclui.
  • Endpoint (opcional) — para stores fora da AWS, como MinIO, Cloudflare R2 ou a interoperabilidade S3 do GCS. Deixe em branco para o AWS S3.
  • Force path-style (opcional) — obrigatório para MinIO/Ceph, precisa estar desligado para AWS S3 e R2.

As chaves de objeto são determinísticas:

{prefix}/{source}/{drainId}/{yyyy}/{mm}/{dd}/{runId}-{seq}.ndjson

Os objetos são gravados com criptografia AES256 do lado do servidor.

Google Cloud Storage

Grava um objeto NDJSON por chunk entregue no seu bucket do GCS.

  • Bucket — o nome do bucket. Precisa já existir; o Studio não cria buckets.
  • Prefix (opcional) — caminho da pasta dentro do bucket. A barra final é opcional.
  • Service account JSON key — cole a chave JSON completa de uma service account com storage.objects.create (e storage.objects.delete, se você quiser que "Test connection" limpe seu objeto de teste). O Studio autentica via JWT OAuth2 de service account e faz o upload pela JSON API do GCS.

Os nomes de objeto seguem o mesmo layout {prefix}/{source}/{drainId}/{yyyy}/{mm}/{dd}/{runId}-{seq}.ndjson do S3. Os metadados do objeto espelham as chaves studio-* do destino S3 por meio de headers x-goog-meta-*.

Azure Blob Storage

Grava um block blob NDJSON por chunk entregue no seu container.

  • Account name — sua storage account (3 a 24 caracteres minúsculos).
  • Container — precisa já existir; o Studio não cria containers.
  • Prefix (opcional) — caminho da pasta dentro do container.
  • Account key — uma chave de acesso da storage account com permissão de escrita no container.

Os nomes de blob seguem o mesmo layout {prefix}/{source}/{drainId}/{yyyy}/{mm}/{dd}/{runId}-{seq}.ndjson. Os metadados studio-* são expostos como metadados de blob do Azure (convertidos para minúsculas conforme as regras de identificador do Azure).

Para nuvens soberanas, defina Endpoint suffix como blob.core.usgovcloudapi.net (US Gov), blob.core.chinacloudapi.cn (China) ou blob.core.cloudapi.de (Alemanha).

Google BigQuery

Envia cada linha em streaming para uma tabela de destino via API tabledata.insertAll, com dedupe por insertId em cada linha.

  • Project ID — seu projeto no GCP (aceita IDs com escopo de domínio, como example.com:my-project).
  • Dataset ID / Table ID — precisam já existir; o Studio não cria tabelas. O schema da tabela precisa acomodar o formato das linhas da fonte (uma coluna por campo de primeiro nível, ou uma única coluna JSON/STRING com o restante como ignoreUnknownValues).
  • Service account JSON key — precisa de roles/bigquery.dataEditor (insert) e roles/bigquery.metadataViewer (para o teste tables.get usado por "Test connection"). O Studio autentica via JWT OAuth2 de service account.

Cada linha é enviada com um insertId no formato {drainId}-{runId}-{sequence}-{index}. O BigQuery deduplica inserts com o mesmo insertId por cerca de 60 segundos, então retentativas dentro dessa janela não duplicam nada. Se um chunk reportar falha parcial (insertErrors), a execução falha indicando os índices das linhas problemáticas, e uma retentativa do driver externo pode duplicar linhas que já tinham sido gravadas — a retentativa limitada do dispatcher minimiza esse risco. Limites aplicados por requisição: 10 MB de corpo, 50.000 linhas.

Snowflake

Insere cada linha em uma coluna VARIANT de destino via Snowflake SQL API v2, com autenticação JWT por par de chaves.

  • Account — o identificador da conta Snowflake. A forma preferida é <orgname>-<acctname> (sem pontos). O formato legado <locator>.<region>.<cloud> também é aceito.
  • User / Warehouse / Database / Schema / Table — precisam já existir. O usuário precisa do privilégio INSERT na tabela e de USAGE no warehouse, no database e no schema.
  • Column (opcional) — nome da coluna VARIANT de destino. O padrão é DATA (compatível com a conversão de identificadores sem aspas do Snowflake).
  • Role (opcional) — role do Snowflake a assumir.
  • Private key (PEM) — chave privada RSA codificada em PKCS8. Registre a chave pública correspondente no usuário do Snowflake com ALTER USER ... SET RSA_PUBLIC_KEY = '...'.

Cada chunk se torna um único INSERT INTO "DB"."SCHEMA"."TABLE" ("col") VALUES (PARSE_JSON(?)), ... com um binding TEXT por linha. Os identificadores são colocados entre aspas para preservar a caixa. O destino trata de forma transparente o padrão assíncrono 202-then-poll do Snowflake. O payload JSON de cada linha é limitado a 16 MB, para corresponder ao limite de VARIANT do Snowflake.

Datadog Logs

Faz POST de cada linha como uma entrada de log na ingestão de logs v2 do Datadog.

  • Site — seu site do Datadog: us1, us3, us5, eu1, ap1, ap2 ou gov.
  • Service (opcional) — valor para o campo reservado service. O padrão é studio.
  • Tags (opcional) — ddtags separadas por vírgula, adicionadas a cada entrada junto com as tags studio_drain_id:, studio_run_id: e studio_source: injetadas automaticamente.
  • API key — uma API key do Datadog (não uma Application key) com permissão de escrita de logs.

Os campos de primeiro nível da linha são indexados automaticamente como atributos de log do Datadog. Os campos reservados ddsource, service, ddtags e message são sempre definidos pelo Studio e sobrescrevem qualquer valor presente na linha. Payloads acima de 1 KB são comprimidos com gzip. Os limites aplicados correspondem aos da ingestão do Datadog: 5 MB por requisição (após compressão), 1000 entradas por requisição, 1 MB por entrada.

Webhook HTTPS

Faz POST de cada chunk como NDJSON no seu endpoint.

  • URL — precisa ser HTTPS. O Studio resolve o hostname e recusa a entrega para IPs privados, de loopback ou de metadados de nuvem. O IP resolvido é fixado durante toda a execução, para impedir DNS rebinding.
  • Signing secret — segredo compartilhado usado para assinatura HMAC-SHA256.
  • Bearer token (opcional) — enviado como Authorization: Bearer <token>.
  • Signature header name (opcional) — o padrão é X-Studio-Signature.

Cada requisição inclui:

Content-Type: application/x-ndjson
User-Agent: Studio-DataDrain/1.0
X-Studio-Timestamp: <unix-seconds>
X-Studio-Signature-Version: v1
X-Studio-Signature: t=<unix-seconds>,v1=<hex(hmac-sha256)>
X-Studio-Drain-Id: <drain id>
X-Studio-Run-Id: <run id>
X-Studio-Source: <source name>
X-Studio-Sequence: <chunk index>
X-Studio-Row-Count: <rows in this chunk>
Idempotency-Key: <runId>-<sequence>

A assinatura é calculada como HMAC-SHA256(secret, "${timestamp}.${body}") e serializada como t=<timestamp>,v1=<hex>. Verifique recalculando sobre a mesma string e rejeitando timestamps com mais de ~5 minutos — isso protege contra ataques de replay com requisições capturadas.

Entregas que falham são repetidas até 3 vezes com backoff exponencial (500ms, 1s, 2s com jitter de ±20%), respeitando Retry-After em 429/503. Respostas 4xx não repetíveis fazem a execução falhar imediatamente.


Agendamento

CadênciaO drain é executado
HourlyUma vez por hora.
DailyUma vez por dia.

Você também pode desativar um drain com o toggle Enabled (ele para de executar, mas é preservado) ou disparar uma execução fora do agendamento com Run now em qualquer linha de drain.


Semântica de entrega

Os drains usam um cursor opaco que só avança em caso de sucesso total. Se uma entrega falhar no meio de uma execução, o cursor permanece inalterado e a próxima execução recomeça da última posição bem-sucedida.

Isso é entrega at-least-once. Combinada com o campo id em cada linha e o header Idempotency-Key em cada chunk de webhook, os sistemas downstream podem deduplicar de forma determinística.

As últimas 10 execuções de cada drain ficam visíveis ao expandir sua linha na página de configurações, com status, contagem de linhas, bytes gravados, localizador do destino (s3://... ou URL do webhook) e a mensagem de erro, caso tenha falhado.


Segurança

  • As credenciais de destino são criptografadas em repouso com a mesma criptografia com suporte a rotação de chaves que protege os tokens OAuth.
  • As credenciais nunca são retornadas pela API do Studio após a criação. As atualizações aceitam novas credenciais; se você omiti-las, o blob criptografado existente é mantido.
  • As URLs de webhook passam por validação contra SSRF: somente HTTPS, sem IPs privados/loopback/metadados, com o IP resolvido fixado para impedir DNS rebinding.
  • Toda chamada de criação, atualização, exclusão, execução manual e teste de conexão é registrada no Log de Auditoria.

Combinando com a Retenção de Dados

Drains e Retenção de Dados são módulos independentes. O Studio não condiciona a retenção ao progresso do drain — se um drain estiver falhando, a retenção continuará expurgando dados no seu próprio agendamento. Esse é o mesmo modelo usado pelo Datadog Archives e pelo AWS CloudWatch + S3 Export: manter as duas configurações ortogonais e deixar que o cliente as combine deliberadamente.

Para usar as duas com segurança, defina a cadência do drain menor que o período de retenção da mesma categoria de dados:

Fonte do drainCombina com a configuração de retenção
Workflow logs, Job logsLog retention
Chats, Chat runsTask cleanup
Audit logs(sem configuração de retenção hoje — logs de auditoria são mantidos indefinidamente)

Por exemplo, com Log retention definido em 30 dias, configure o drain de workflow logs como Hourly ou Daily, para que cada linha seja exportada bem antes de a retenção expurgá-la do Studio. Acompanhe as execuções recentes do drain na página de configurações; se um drain estiver falhando por mais tempo que sua janela de retenção, você pode perder linhas que a retenção expurga antes da exportação.

Depois que os dados chegam ao seu bucket ou sistema de webhook, o ciclo de vida do arquivamento (transições para o Glacier, expiração, propagação do direito ao esquecimento do GDPR) é governado pela sua própria infraestrutura — o Studio não tem mais visibilidade sobre esses dados depois que a entrega é concluída com sucesso.


Common Questions

Apenas owners e admins da organização podem ver, criar, editar, executar ou excluir drains. No Studio Cloud, a organização precisa estar em um plano Enterprise.
O cursor do drain só avança em caso de sucesso geral, então uma falha faz a próxima execução reprocessar os mesmos chunks. Cada linha tem um campo `id` estável e cada chunk de webhook tem um header `Idempotency-Key`, para que os receptores possam deduplicar.
Sim — crie um drain por fonte, todos apontando para o mesmo bucket ou endpoint. Destinos S3 separam por fonte automaticamente; receptores de webhook podem ramificar pelo header `X-Studio-Source`.
Não. A exclusão remove apenas a configuração do drain e seu histórico de execuções do Studio. Os dados já gravados no seu bucket ou enviados ao seu webhook são seus e não são afetados.
A execução falha, o cursor do drain não avança e a execução com falha é registrada com o erro. Depois que você corrigir as credenciais com um Update ou recriando o drain, a próxima execução recomeça de onde a última execução bem-sucedida parou.
NDJSON — JSON delimitado por quebras de linha, uma linha por registro. Cada chunk é um único objeto no S3 ou um único corpo de POST.

Configuração em auto-hospedagem

Variáveis de ambiente

DATA_DRAINS_ENABLED=true
NEXT_PUBLIC_DATA_DRAINS_ENABLED=true

NEXT_PUBLIC_DATA_DRAINS_ENABLED exibe a página Settings → Enterprise → Data Drains na interface. DATA_DRAINS_ENABLED controla os endpoints de mutação no servidor e o dispatcher do cron — quando não está definido em uma instalação auto-hospedada, requisições de criação/atualização/exclusão/execução de drains retornam 404 e o dispatcher não faz nada. As duas devem ser definidas como true juntas.

Fora disso, os Data Drains usam a mesma infraestrutura de jobs em segundo plano do Trigger.dev empregada no restante do Studio — nenhuma configuração adicional é necessária. O dispatcher do cron roda a cada hora e distribui os drains pendentes como jobs em segundo plano.