Pular para o conteúdo principal

Conectar ao Amazon Kinesis

Use o Structured Streaming para ler e gravar dados no Amazon Kinesis.

A Databricks recomenda que você habilite Endpoint de VPC do S3 para que todo o tráfego do S3 seja roteado na rede da AWS.

nota

Você pode refragmentar com o Structured Streaming aumentando o número de fragmentos sem interromper ou reiniciar a transmissão.

Para recomendações sobre a solução de problemas de latência de query, consulte Recomendações para reduzir a latência com o Kinesis.

Autenticação

O Kinesis oferece suporte à autenticação com uma conexão do Unity Catalog, uma credencial de serviço ou métodos alternativos, como um instance profile ou chaves de acesso. Consulte Autenticação.

Esquema

O Kinesis retorna registros com o seguinte esquema:

Coluna

Tipo

Descrição

partitionKey

string

A chave de partição que identifica a qual fragmento o registro está atribuído.

data

binário

O blob de dados do registro, como binário opaco.

stream

string

O nome ou ARN da transmissão do Kinesis da qual o registro foi lido.

shardId

string

O ID do fragmento do qual o registro foi lido.

sequenceNumber

string

O identificador exclusivo do registro dentro de seu fragmento.

approximateArrivalTimestamp

carimbo de data/hora

O horário aproximado em que o registro foi inserido na transmissão.

Coluna

Tipo

Descrição

partitionKey

string

A chave de partição que identifica a qual fragmento o registro está atribuído.

data

binário

O blob de dados do registro, como binário opaco.

stream

string

O nome ou ARN da transmissão do Kinesis da qual o registro foi lido.

shardId

string

O ID do fragmento do qual o registro foi lido.

sequenceNumber

string

O identificador exclusivo do registro dentro de seu fragmento.

approximateArrivalTimestamp

carimbo de data/hora

O horário aproximado em que o registro foi inserido na transmissão.

Para desserializar os dados na coluna data, converta o campo para uma string.

Início rápido

O notebook a seguir demonstra como executar o WordCount com a Transmissão estruturada com Kinesis.

Notebook Kinesis WordCount com Structured Streaming

Configurar opções do Kinesis

No Databricks Runtime 13.3 LTS e superior, você pode usar Trigger.AvailableNow com o Kinesis. Consulte Ingerir registros do Kinesis como um lote incremental.

No Databricks Runtime 16.1 e superior, você pode usar streamARN para identificar fontes do Kinesis. Para todas as versões do Databricks Runtime, você deve especificar streamName ou streamARN, mas não ambos.

atenção

Não alterne entre streamName e streamARN para uma query de transmissão ativa. O Databricks não oferece suporte à alternância entre essas opções durante a transmissão. Reiniciar a query pode resultar em registros duplicados ou perda de dados. Para alternar de streamName para streamARN, inicie uma nova query de transmissão com um diretório de ponto de verificação novo.

Para obter a lista completa de opções, consulte Kinesis.

Adicionar ou remover fontes de transmissão

No Databricks Runtime 19 e acima, você pode alterar as streams de origem do Kinesis para queries de Structured Streaming usando opções do Spark, streamName ou streamARN.

Adicionar uma transmissão

Para adicionar uma transmissão, inclua-a na lista de opções streamName ou streamARN e reinicie a transmissão. Para cada nova fonte de transmissão, a query lê o deslocamento disponível mais antigo dos shards da fonte.

O exemplo a seguir usa a opção streamName para adicionar a fonte do Kinesis stream3 a uma query que lia anteriormente stream1 e stream2:

Python
df = (spark.readStream
.format("kinesis")
.option("streamName", "stream1,stream2")
.load()
)

df.stop()

df = (spark.readStream
.format("kinesis")
.option("streamName", "stream1,stream2,stream3") # Previous value was "stream1,stream2"
.load()
)

Remover uma transmissão

Por default, quando você remove uma fonte de transmissão do Kinesis da lista de opções streamName ou streamARN, a query falha na reinicialização com um erro KINESIS_SOURCE_STREAMS_REMOVED_ON_RESTART. Isso garante que a query não ignore silenciosamente registros não lidos da transmissão removida.

Para remover uma fonte de transmissão do Kinesis, faça o seguinte:

  1. Defina spark.databricks.kinesis.failOnDataLoss como false na configuração do Spark do cluster e reinicie o cluster. Para obter mais informações sobre failOnDataLoss, consulte Lidar com perda de dados.

  2. Remova a transmissão da opção streamName ou streamARN e reinicie a query. Por exemplo, para parar de ler stream2:

    Python
    df = (spark.readStream
    .format("kinesis")
    .option("streamName", "stream1,stream2")
    .load()
    )

    df.stop()

    df = (spark.readStream
    .format("kinesis")
    .option("streamName", "stream1") # Previous value was "stream1,stream2"
    .load()
    )

Você só precisa definir spark.databricks.kinesis.failOnDataLoss como false e reiniciar o cluster uma vez. Depois disso, a remoção de fontes de transmissão adicionais nesse cluster requer apenas uma reinicialização da query, não outra reinicialização do cluster.

nota

Por default, a remoção de uma fonte de transmissão do Kinesis não cancela o registro do consumidor de distribuição avançada (EFO) da fonte de transmissão, o que pode continuar a gerar custos do provedor de cloud. Para cancelar o registro do consumidor quando a query parar, defina a opção requireConsumerDeregistration como true. See Kinesis.

Para gerenciar consumidores diretamente, consulte Configurar a distribuição avançada (EFO) do Kinesis para leituras de query de transmissão.

Monitoramento e alerta de baixa latência

Casos de uso de alerta exigem baixa latência. Para minimizar a latência:

  • Verifique se sua query de transmissão é a única consumidora da transmissão do Kinesis para otimizar o desempenho da busca e evitar os limites de taxa do Kinesis.
  • Defina a opção maxFetchDuration com um valor pequeno, como 200 ms, para processar os dados obtidos o mais rápido possível. Esta opção é um compromisso: ela prioriza uma velocidade de processamento mais rápida por lote em vez de uma garantia de que os registros mais recentes sejam consumidos em cada lote. Por exemplo, se você usar Trigger.AvailableNow, um valor pequeno pode fazer com que sua query fique atrasada em relação aos registros mais recentes na transmissão do Kinesis.
  • Defina a opção minFetchPeriod como 210ms para buscar com a maior frequência possível.
  • Configure a opção shardsPerTask ou configure o cluster de forma que # cores in cluster >= 2 * (# Kinesis shards) / shardsPerTask. Isso garante que as tarefas de pré-busca em segundo plano e as tarefas de query de transmissão sejam executadas simultaneamente.

Se sua query estiver recebendo dados a cada 5 segundos, você poderá exceder os limites de taxa do Kinesis. Revise suas configurações.

Monitorar métricas do Kinesis

O Kinesis relata o número de milissegundos em que um consumidor está atrasado em relação ao início de uma transmissão para cada workspace. As métricas avgMsBehindLatest, maxMsBehindLatest e minMsBehindLatest fornecem a média, o mínimo e o máximo de milissegundos em todos os Workspace no processo de query de transmissão. Consulte Monitoramento de queries de Structured Streaming no Databricks.

Se você estiver executando o stream em um notebook, veja as métricas na tab Dados brutos no painel de progresso da query de transmissão. Aqui está um exemplo:

JSON
{
"sources": [
{
"description": "KinesisV2[stream]",
"metrics": {
"avgMsBehindLatest": "32000.0",
"maxMsBehindLatest": "32000",
"minMsBehindLatest": "32000"
}
}
]
}

Ingerir registros do Kinesis como um lote incremental

No Databricks Runtime 13.3 LTS e acima, o Databricks oferece suporte ao uso de Trigger.AvailableNow com fontes de dados Kinesis para semântica de lote incremental. O seguinte descreve a configuração básica:

  1. Quando um micro-batch lê gatilhos no modo disponível agora, a hora atual é registrada pelo cliente Databricks.
  2. O Databricks pesquisa o sistema de origem para todos os registros com carimbos de data/hora entre esse tempo registrado e o ponto de verificação anterior.
  3. O Databricks carrega esses registros com Trigger.AvailableNow semântica.

O Databricks usa um mecanismo de melhor esforço para tentar consumir todos os registros que existem nas transmissões do Kinesis quando a query de transmissão é executada. Devido a pequenas possíveis diferenças nos carimbos de data e à falta de garantia de ordenação nas fontes de dados, um lote acionado pode não incluir alguns registros. Os registros omitidos são processados no próximo micro-Trigger acionado.

nota

Se a query continuar falhando ao buscar registros da transmissão do Kinesis mesmo havendo registros, tente aumentar o valor de maxFetchDuration.

Consulte AvailableNow: Processamento de lotes incrementais.

Lidar com perda de dados

atenção

Use failOnDataLoss apenas se sua carga de trabalho puder tolerar registros ausentes. O uso incorreto disso pode resultar em perda permanente de dados. Se você não puder tolerar registros ausentes, reinicie a transmissão com um novo ponto de verificação para reprocessar todos os registros.

O Databricks recomenda que você use isso apenas como uma mitigação temporária para um problema de perda de dados. Investigue e corrija a causa raiz, como um período de retenção do Kinesis que seja muito curto.

Se os registros de um shard do Kinesis expirarem antes que sua query de transmissão os leia, ou se você excluir e recriar uma transmissão do Kinesis com o mesmo nome, a query falhará com um erro KINESIS_COULD_NOT_READ_SHARD_UNTIL_END_OFFSET. Veja KINESIS_COULD_NOT_READ_SHARD_UNTIL_END_OFFSET.

Por default, as consultas de transmissão falham quando detectam possível perda de dados. Para configurar a query para ignorar registros ilegíveis e continuar o processamento, defina spark.databricks.kinesis.failOnDataLoss como false na configuração do Spark do cluster e reinicie o cluster.

Gravar no Kinesis

Use o seguinte trecho de código como um ForeachSink para gravar dados no Kinesis. Exige um Dataset[(String, Array[Byte])].

nota

O trecho de código a seguir fornece semântica pelo menos uma vez , não exatamente uma vez.

Notebook Kinesis Foreach Sink

Recomendações para reduzir a latência com o Kinesis

Esta seção contém recomendações para a solução de problemas de várias causas de latência para transmissões do Kinesis.

A origem do Kinesis executa jobs do Spark em um thread em segundo plano para buscar previamente e periodicamente os dados do Kinesis e, em seguida, armazená-los em cache na memória dos executores do Spark. Após a conclusão de cada passo de pré-busca, a query de transmissão pode processar os dados em cache. O passo de pré-busca afeta significativamente a latência de ponta a ponta observada e o throughput.

Reduzir a latência de pré-busca

Para otimizar a latência mínima de query e o uso máximo de recursos, use o seguinte cálculo:

total number of CPU cores in the cluster (across all executors) >= total number of Kinesis shards / shardsPerTask.

importante

minFetchPeriod pode criar várias chamadas de API GetRecords para o fragmento do Kinesis até atingir ReadProvisionedThroughputExceeded. Se ocorrer uma exceção, pode não ser um problema, pois o conector maximiza a utilização do fragmento do Kinesis.

Evite lentidões causadas por excesso de erros de limite de taxa

O conector reduz pela metade a quantidade de dados lidos do Kinesis cada vez que encontra um erro de limitação de taxa e registra esse evento no log com uma mensagem: "Hit rate limit. Sleeping for 5 seconds."

Você poderá ver esses erros enquanto uma transmissão está sendo atualizada. Se você vir esses erros após uma transmissão ser atualizada, talvez precise ajustar a carga de trabalho aumentando a capacidade do Kinesis na AWS ou ajustando as opções de pré-busca no Spark.

Evite o spill de disco

Se você tiver um aumento repentino no volume de dados em suas transmissões do Kinesis, a capacidade de buffer atribuída poderá ser preenchida e não esvaziar rápido o suficiente para adicionar novos dados. O Spark faz spill de dados do buffer para o disco, o que reduz a velocidade do processamento de transmissão, e um evento aparece no log com uma mensagem como a seguinte:

Bash
./log4j.txt:879546:20/03/02 17:15:04 INFO BlockManagerInfo: Updated kinesis_49290928_1_ef24cc00-abda-4acd-bb73-cb135aed175c on disk on 10.0.208.13:43458 (current size: 88.4 MB, original size: 0.0 B)

Para evitar spill, aumente a capacidade de memória do cluster adicionando mais nós ou aumentando a memória por nó, ou reduza o parâmetro de configuração fetchBufferSize.

Tarefas de gravação no S3 suspensas

Habilite a especulação do Spark para encerrar tarefas suspensas que impediriam o processamento da transmissão de prosseguir. Para garantir que as tarefas não sejam encerradas de forma muito agressiva, ajuste cuidadosamente o quantil e o multiplicador para essa configuração. A Databricks recomenda que você configure spark.speculation.multiplier para 3 e spark.speculation.quantile para 0.95 e ajuste conforme necessário.

Reduza a latência do ponto de verificação em transmissões com estado

O Databricks recomenda o uso do RocksDB com ponto de verificação do log de alterações para consultas de transmissão com estado. Consulte Habilitar ponto de verificação do log de alterações.

Configurar o fan-out aprimorado (EFO) do Kinesis para leituras de query de transmissão

No Databricks Runtime 11.3e acima, o conector Databricks Runtime Kinesis fornece suporte para usar o recurso de fan-out aprimorado (EFO) do Amazon Kinesis.

O fan-out aprimorado do Kinesis fornece throughput dedicado de 2 MB/s por fragmento por consumidor (máximo de 20 consumidores por transmissão) e entrega registros no modo push em vez do modo pull.

Por default, uma query de Structured Streaming configurada com o modo EFO registra-se como um consumidor com throughput dedicado e um nome de consumidor e ARN (Amazon Resource Name) de consumidor exclusivos no Kinesis Data Streams.

Por default, o Databricks usa o ID da query de transmissão com o prefixo databricks_ para nomear o novo consumidor. Opcionalmente, você pode especificar as opções consumerNamePrefix ou consumerName para substituir este comportamento. O consumerName deve ser uma string que contenha apenas letras, números e os caracteres especiais _ . -.

Ao reiniciar a query, a fonte do Kinesis usa o modo de polling para reproduzir o último lote não confirmado, caso exista. Após a transmissão reproduzir o lote não confirmado, a fonte alterna de volta para o modo EFO para leituras subsequentes.

importante

Um consumidor registrado do EFO incorre em taxas adicionais na Amazon Kinesis. Para cancelar o registro do consumidor automaticamente na desativação da query, defina a opção requireConsumerDeregistration como true. O Databricks não pode garantir o cancelamento do registro em eventos como falhas no driver ou falhas nos nós. Em caso de falha no job, a Databricks recomenda gerenciar os consumidores registrados diretamente para evitar cobranças excessivas do Kinesis.

Gerenciamento de consumidores off-line usando um notebook Databricks

Use o utilitário AWSKinesisConsumerManager para registrar, listar ou cancelar o registro de consumidores para fluxos de dados do Kinesis programaticamente, em vez de configurar consumidores manualmente no console da sua account AWS. Por exemplo, use o utilitário para criar um consumidor para uma nova transmissão ou, se você planeja interromper uma transmissão permanentemente, use o utilitário para excluir o consumidor na AWS.

A utilidade de gerenciador de consumidor está disponível apenas em Scala com o compute definido para o modo de acesso dedicado. Consulte Modos de acesso.

Para usar este utilitário em um notebook Databricks:

  1. Em um novo notebook do Databricks anexado a um cluster ativo, crie um AWSKinesisConsumerManager com as informações de autenticação necessárias.

    Scala
    import com.databricks.sql.kinesis.AWSKinesisConsumerManager

    val manager = AWSKinesisConsumerManager.newManager()
    .option("serviceCredential", serviceCredentialName)
    .option("region", kinesisRegion)
    .create()
  2. Liste e exiba os consumidores.

    Scala
    val consumers = manager.listConsumers("<stream name>")
    display(consumers)
  3. Registre um consumidor para a transmissão informada.

    Scala
    val consumerARN = manager.registerConsumer("<stream name>", "<consumer name>")
  4. Cancele o registro de um consumidor para a transmissão informada.

    Scala
    manager.deregisterConsumer("<stream name>", "<consumer name>")