Aller au contenu principal

Connectez-vous à Amazon Kinesis

Utilisez Structured Streaming pour lire et écrire des données vers Amazon Kinesis.

Databricks recommande d'activer les endpoints S3 Virtual Private Cloud (VPC) afin que tout le trafic S3 soit routé sur le réseau AWS.

remarque

Vous pouvez repartitionner avec Structured Streaming en augmentant le nombre de partitions sans interrompre ni redémarrer le Stream.

Pour obtenir des recommandations sur le dépannage de la latence des requêtes, consultez Recommendations for reducing latency with Kinesis.

Authentification

Kinesis prend en charge l’authentification avec une connexion Unity Catalog, un identifiant de service ou d’autres méthodes telles qu’un profil d’instance ou des clés d’accès. Consultez Authentification.

Schéma

Kinesis renvoie des enregistrements avec le schéma suivant :

Colonne

Type

Description

partitionKey

chaîne

La clé de partition qui identifie la partition à laquelle l’enregistrement est affecté.

data

binaire

Le blob de données de l’enregistrement, sous forme binaire opaque.

stream

chaîne

Le nom ou l’ARN du Stream Kinesis à partir duquel l’enregistrement a été lu.

shardId

chaîne

L’ID de la partition à partir de laquelle l’enregistrement a été lu.

sequenceNumber

chaîne

L’identifiant unique de l’enregistrement au sein de sa partition.

approximateArrivalTimestamp

Horodatage

L’heure approximative à laquelle l’enregistrement a été inséré dans le Stream.

Colonne

Type

Description

partitionKey

chaîne

La clé de partition qui identifie la partition à laquelle l’enregistrement est affecté.

data

binaire

Le blob de données de l’enregistrement, sous forme binaire opaque.

stream

chaîne

Le nom ou l’ARN du Stream Kinesis à partir duquel l’enregistrement a été lu.

shardId

chaîne

L’ID de la partition à partir de laquelle l’enregistrement a été lu.

sequenceNumber

chaîne

L’identifiant unique de l’enregistrement au sein de sa partition.

approximateArrivalTimestamp

Horodatage

L’heure approximative à laquelle l’enregistrement a été inséré dans le Stream.

Pour désérialiser les données de la colonne data, convertissez le champ en chaîne (string).

Guide de démarrage rapide

Le notebook suivant montre comment exécuter WordCount en utilisant Structured Streaming avec Kinesis.

Notebook Kinesis WordCount avec Structured Streaming

Configurer les options Kinesis

Dans Databricks Runtime 13.3 LTS et versions ultérieures, vous pouvez utiliser Trigger.AvailableNow avec Kinesis. Consultez Ingérer des enregistrements Kinesis en tant que batch incrémentiel.

Dans Databricks Runtime 16.1 et versions supérieures, vous pouvez utiliser streamARN pour identifier les sources Kinesis. Pour toutes les versions de Databricks Runtime, vous devez spécifier soit streamName, soit streamARN, mais pas les deux.

attention

Ne basculez pas entre streamName et streamARN pour une query de streaming active. Databricks ne prend pas en charge le basculement entre ces options en cours de Stream. Le redémarrage de la query peut entraîner des enregistrements en double ou une perte de données. Pour passer de streamName à streamARN, start une nouvelle query de streaming avec un répertoire de point de contrôle vierge.

Pour la liste complète des options, consultez Kinesis.

Ajouter ou supprimer des sources de Stream

Dans Databricks Runtime 19 et versions ultérieures, vous pouvez modifier les flux sources Kinesis pour les requêtes Structured Streaming à l'aide des options Spark, streamName ou streamARN.

Ajouter un Stream

Pour ajouter un Stream, incluez-le dans la liste d'options streamName ou streamARN et redémarrez le Stream. Pour chaque nouvelle source de Stream, la query lit l'offset disponible le plus ancien à partir des shards de la source.

L'exemple suivant utilise l'option streamName pour ajouter la source Kinesis stream3 à une query qui lisait précédemment stream1 et 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()
)

Supprimer un Stream

Par défaut, lorsque vous supprimez une source de Stream Kinesis de la liste d’options streamName ou streamARN, la query échoue au redémarrage avec une erreur KINESIS_SOURCE_STREAMS_REMOVED_ON_RESTART. Cela garantit que la query ne saute pas silencieusement les enregistrements non lus du Stream supprimé.

Pour supprimer une source de Stream Kinesis, procédez comme suit :

  1. Définissez spark.databricks.kinesis.failOnDataLoss sur false dans la configuration Spark du cluster et redémarrez le cluster. Pour plus d’informations sur failOnDataLoss, consultez Handle data loss.

  2. Supprimez le Stream de l’option streamName ou streamARN et redémarrez la query. Par exemple, pour arrêter la lecture de 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()
    )

Il vous suffit de définir spark.databricks.kinesis.failOnDataLoss sur false et de redémarrer le cluster une seule fois. Ensuite, la suppression de sources de Stream supplémentaires sur ce cluster ne nécessite qu'un redémarrage de la query, et non un nouveau redémarrage du cluster.

remarque

Par default, la suppression d'une source de Stream Kinesis ne désinscrit pas le consommateur EFO (enhanced fan-out) de la source de Stream, ce qui pourrait continuer à engendrer des coûts auprès du fournisseur cloud. Pour désinscrire le consommateur lorsque la query s'arrête, définissez l'option requireConsumerDeregistration sur true. See Kinesis.

Pour gérer les consommateurs directement, consultez Configurer Kinesis enhanced fan-out (EFO) pour les lectures de query en streaming.

Monitoring et alertes à faible latence

Les cas d'usage d'alerte nécessitent une faible latence. Pour minimiser la latence :

  • Vérifiez que votre query de streaming est le seul consommateur du stream Kinesis afin d'optimiser les performances de récupération et d'éviter les limites de débit Kinesis.
  • Définissez l’option maxFetchDuration sur une petite valeur, telle que 200 ms, pour traiter les données récupérées aussi rapidement que possible. Cette option est un compromis : elle privilégie une vitesse de traitement plus rapide par batch plutôt qu’une garantie que les enregistrements les plus récents sont consommés dans chaque batch. Par exemple, si vous utilisez Trigger.AvailableNow, une petite valeur pourrait entraîner un retard de votre query par rapport aux enregistrements les plus récents dans le Stream Kinesis.
  • Définissez l'option minFetchPeriod sur 210 ms pour effectuer une récupération aussi fréquemment que possible.
  • Définissez l'option shardsPerTask ou configurez le cluster de sorte que # cores in cluster >= 2 * (# Kinesis shards) / shardsPerTask. Cela garantit que les tâches de préchargement en arrière-plan et les tâches de query de streaming s'exécutent simultanément.

Si votre query reçoit des données toutes les 5 secondes, vous risquez de dépasser les limites de débit Kinesis. Vérifiez vos configurations.

Surveiller les indicateurs Kinesis

Kinesis indique le nombre de millisecondes de retard d'un consommateur par rapport au début d'un stream pour chaque workspace. Les métriques avgMsBehindLatest, maxMsBehindLatest et minMsBehindLatest fournissent la moyenne, le minimum et le maximum de millisecondes sur tous les Workspace dans le processus de streaming query. Consultez Monitoring des requêtes Structured Streaming sur Databricks.

Si vous exécutez le stream dans un notebook, consultez les métriques sous l'onglet Données brutes dans le tableau de bord de progression de la query de streaming. Voici un exemple :

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

Ingérer les enregistrements Kinesis sous forme de batch incrémentiel

Dans Databricks Runtime 13.3 LTS et versions ultérieures, Databricks prend en charge l’utilisation de Trigger.AvailableNow avec des sources de données Kinesis pour une sémantique de batch incrémentielle. La section suivante décrit la configuration de base :

  1. Lorsqu'une lecture micro-batch se Trigger en mode « available now », l'heure actuelle est enregistrée par le client Databricks.
  2. Databricks interroge le système source pour obtenir tous les enregistrements dont les Timestamp sont compris entre ce moment enregistré et le point de contrôle précédent.
  3. Databricks charge ces enregistrements en utilisant la sémantique Trigger.AvailableNow.

Databricks utilise un mécanisme « best-effort » pour tenter de consommer tous les enregistrements existant dans les flux Kinesis lors de l'exécution de la streaming query. En raison de légères différences potentielles dans les Timestamp et d'un manque de garantie sur l'ordre dans les sources de données, un Triggered batch peut ne pas inclure certains enregistrements. Les enregistrements omis sont traités lors du micro-batch déclenché suivant.

remarque

Si la requête continue d’échouer à récupérer des enregistrements du Kinesis Stream alors qu’il en existe, essayez d’augmenter la valeur maxFetchDuration.

Voir AvailableNow: Traitement incrémentiel par batch.

Gérer la perte de données

attention

Utilisez failOnDataLoss uniquement si votre workload peut tolérer des enregistrements manquants. Une utilisation incorrecte peut entraîner une perte de données permanente. Si vous ne pouvez pas tolérer d’enregistrements manquants, redémarrez le stream avec un nouveau point de contrôle pour retraiter tous les enregistrements.

Databricks vous recommande de n’utiliser ceci que comme mesure d’atténuation temporaire pour un problème de perte de données. Enquêtez sur la cause principale et corrigez-la, par exemple une période de rétention Kinesis trop courte.

Si les enregistrements d'un shard Kinesis expirent avant que votre query de streaming ne les lise, ou si vous supprimez et recréez un stream Kinesis avec le même nom, la query échoue avec une erreur KINESIS_COULD_NOT_READ_SHARD_UNTIL_END_OFFSET. Voir KINESIS_COULD_NOT_READ_SHARD_UNTIL_END_OFFSET.

Par default, les requêtes en streaming échouent lorsqu'elles détectent une perte de données potentielle. Pour configurer la requête afin d'ignorer les enregistrements illisibles et poursuivre le traitement, définissez spark.databricks.kinesis.failOnDataLoss sur false dans la configuration Spark du cluster et redémarrez le cluster.

Écrire vers Kinesis

Utilisez l'extrait de code suivant comme ForeachSink pour écrire des données dans Kinesis. Cela nécessite un Dataset[(String, Array[Byte])].

remarque

L'extrait de code suivant fournit une sémantique at least once (au moins une fois), et non exactly once (exactement une fois).

Notebook Kinesis Foreach Sink

Recommandations pour réduire la latence avec Kinesis

Cette section contient des recommandations pour résoudre les diverses causes de latence des flux Kinesis.

La source Kinesis exécute des Jobs Spark dans un thread d'arrière-plan pour précharger périodiquement les données Kinesis, puis les mettre en cache dans la mémoire de l'exécuteur Spark. Une fois chaque étape de préchargement terminée, la query de streaming peut traiter les données mises en cache. L'étape de préchargement affecte considérablement la latence et le throughput de bout en bout observés.

Réduire la latence de préchargement

Pour optimiser la latence des query et maximiser l’utilisation des ressources, utilisez le calcul suivant :

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

important

minFetchPeriod peut créer plusieurs appels d’API GetRecords vers la partition Kinesis jusqu’à ce qu’elle atteigne ReadProvisionedThroughputExceeded. Si une exception se produit, ce n’est peut-être pas un problème, car le connecteur maximise l’utilisation de la partition Kinesis.

Évitez les ralentissements causés par un trop grand nombre d’erreurs de limite de débit

Le connecteur réduit de moitié la quantité de données lues depuis Kinesis chaque fois qu’il rencontre une erreur de limitation de débit et enregistre cet événement dans le log avec un message : "Hit rate limit. Sleeping for 5 seconds."

Vous pourriez rencontrer ces erreurs pendant qu’un stream rattrape son retard. Si vous voyez ces erreurs après qu’un stream a rattrapé son retard, vous devrez peut-être ajuster la charge de travail en augmentant la capacité Kinesis dans AWS ou en ajustant les options de prélecture dans Spark.

Éviter le spill sur disque

Si vous constatez une augmentation soudaine du volume de données dans vos Kinesis Stream, la capacité de tampon allouée pourrait se remplir et ne pas se vider assez rapidement pour ajouter de nouvelles données. Spark **spill** les données du tampon vers le disque, ce qui ralentit le **Stream processing**, et un événement apparaît dans les **logs** avec un message tel que le suivant :

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)

Pour éviter le spill, augmentez la capacité mémoire du cluster en ajoutant des nœuds ou en augmentant la mémoire par nœud, ou réduisez le parameter de configuration fetchBufferSize.

Tâches d'écriture S3 suspendues

Activez la spéculation Spark pour terminer les tâches suspendues qui empêcheraient le traitement du stream de se poursuivre. Pour vous assurer que les tâches ne sont pas terminées de manière trop agressive, ajustez soigneusement le quantile et le multiplicateur pour ce paramètre. Databricks vous recommande de définir spark.speculation.multiplier sur 3 et spark.speculation.quantile sur 0.95, puis d'ajuster selon vos besoins.

Réduire la latence liée à la création de points de contrôle dans les flux avec état

Databricks recommande d’utiliser RocksDB avec le point de contrôle de journal des modifications pour les queries de streaming avec état. Consultez Activer le point de contrôle de journal des modifications.

Configurer le fan-out amélioré (EFO) de Kinesis pour les lectures de streaming query

Dans Databricks Runtime 11.3 et versions ultérieures, le connecteur Databricks Runtime Kinesis prend en charge l’utilisation de la fonctionnalité Amazon Kinesis enhanced fan-out (EFO).

Kinesis enhanced fan-out fournit un throughput dédié de 2 Mo/s par shard et par consommateur (maximum de 20 consommateurs par stream), et délivre les enregistrements en mode push plutôt qu’en mode pull.

Par default, une query Structured Streaming configurée avec le mode EFO s'enregistre en tant que consommateur avec un throughput dédié, un nom de consommateur unique et un ARN (Amazon Resource Name) de consommateur dans Kinesis Data Streams.

Par default, Databricks utilise l'ID de query de streaming avec le préfixe databricks_ pour nommer le nouveau consommateur. Vous pouvez éventuellement spécifier les options consumerNamePrefix ou consumerName pour ignorer ce comportement. Le consumerName doit être une chaîne contenant uniquement des lettres, des chiffres et les caractères spéciaux _ . -.

Au redémarrage de la query, la source Kinesis utilise le mode d'interrogation pour rejouer le dernier batch non validé, s'il en existe un. Une fois que le Stream a rejoué le batch non validé, la source repasse en mode EFO pour les lectures ultérieures.

important

Un consommateur EFO enregistré entraîne des frais supplémentaires sur Amazon Kinesis. Pour désinscrire automatiquement le consommateur lors du démontage de la query, définissez l'option requireConsumerDeregistration sur true. Databricks ne peut pas garantir la désinscription en cas d'événements tels que des plantages du driver ou des défaillances de nœuds. En cas d'échec du Job, Databricks recommande de gérer directement les consommateurs enregistrés afin d'éviter des frais Kinesis excessifs.

Gestion des consommateurs hors ligne à l’aide d’un notebook Databricks

Utilisez l’utilitaire AWSKinesisConsumerManager pour enregistrer, lister ou désenregistrer par programmation des consommateurs pour les flux de données Kinesis, au lieu de configurer manuellement les consommateurs dans la console de votre compte AWS. Par exemple, utilisez l’utilitaire pour créer un consommateur pour un nouveau Stream ou, si vous prévoyez d’arrêter définitivement un Stream, utilisez l’utilitaire pour supprimer le consommateur dans AWS.

L’utilitaire de gestion des consommateurs n’est disponible qu’en Scala avec le compute configuré en mode d’accès dédié. Voir Modes d’accès.

Pour utiliser cette infrastructure publique dans un Notebook Databricks :

  1. Dans un nouveau Notebook Databricks associé à un cluster actif, créez un AWSKinesisConsumerManager avec les informations d’authentification requises.

    Scala
    import com.databricks.sql.kinesis.AWSKinesisConsumerManager

    val manager = AWSKinesisConsumerManager.newManager()
    .option("serviceCredential", serviceCredentialName)
    .option("region", kinesisRegion)
    .create()
  2. Lister et afficher les consommateurs.

    Scala
    val consumers = manager.listConsumers("<stream name>")
    display(consumers)
  3. Enregistrez un consommateur pour le stream donné.

    Scala
    val consumerARN = manager.registerConsumer("<stream name>", "<consumer name>")
  4. Désenregistrez un consommateur pour le stream donné.

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