Connectez-vous à Amazon Kinesis
Utilisez Structured Streaming pour lire et écrire des données vers Amazon Kinesis.
Databricks vous recommande d’activer les Endpoint S3 Virtual Private Cloud (VPC) afin que tout le trafic S3 soit acheminé sur le réseau AWS.
Si vous supprimez et recréez un Stream Kinesis, vous ne pouvez pas réutiliser les répertoires de point de contrôle existants pour redémarrer une query de streaming. Vous devez supprimer les répertoires de point de contrôle et start ces queries à partir de zéro. Vous pouvez repartitionner avec Structured Streaming en augmentant le nombre de shards sans interrompre ni redémarrer le Stream.
Pour des recommandations sur le dépannage de la latence des requêtes, consultez Recommandations pour réduire la latence avec Kinesis.
Authentification avec Amazon Kinesis
Dans Databricks Runtime 16.1 et versions ultérieures, Databricks vous recommande de gérer les connexions à Kinesis avec un identifiant de service Databricks. Voir Créer des identifiants de service.
Pour utiliser un identifiant de service, procédez comme suit :
- Créez une information d'identification de service Databricks à l'aide d'un rôle IAM avec les autorisations nécessaires pour accéder à Kinesis. Voir Étape 1 : Créez un rôle IAM.
- Fournissez le nom de l'identifiant de service à l'aide de l'option
serviceCredentiallors de la définition d'une lecture en streaming.
La source Kinesis nécessite les autorisations ListShards, GetRecords, et GetShardIterator. Si vous rencontrez Amazon: Access Denied exceptions, vérifiez que votre rôle IAM dispose de ces autorisations. Consultez Contrôler l'accès aux Ressources Amazon Kinesis Data Stream à l'aide d'IAM.
Autres méthodes d'authentification
Dans Databricks Runtime 16,0 et versions inférieures, les informations d'identification de service Databricks ne sont pas disponibles. Databricks propose les méthodes d'authentification alternatives suivantes :
Profil d'instance
Attacher un profil d'instance lors de la configuration du compute. Consultez les profils d'instance.
Les profils d'instance ne sont pas pris en charge en mode d'accès standard (anciennement mode d'accès partagé). Consultez Exigences et limitations du compute standard.
Utiliser les clés d'accès directement
Définissez les options awsAccessKey et awsSecretKey.
Si vous utilisez des clés, stockez-les à l'aide des secrets Databricks. Consultez Gestion des secrets.
Assumez le rôle IAM
Certaines configurations de compute vous permettent d'assumer un rôle IAM à l'aide de l'option roleArn. Pour assumer un rôle, lancez votre cluster avec des autorisations pour assumer le rôle ou fournissez des clés d'accès via awsAccessKey et awsSecretKey.
Cette méthode prend en charge l'authentification inter-comptes. Pour plus d'informations, consultez Déléguer l'accès entre les comptes AWS à l'aide de rôles IAM.
Vous pouvez éventuellement spécifier l’ID externe avec roleExternalId et un nom de session avec roleSessionName.
Schéma
Kinesis renvoie des enregistrements avec le schéma suivant :
Colonne | Type |
|---|---|
| chaîne |
| binaire |
| chaîne |
| chaîne |
| chaîne |
| Horodatage |
Pour désérialiser les données dans la colonne data, convertissez le champ en chaîne.
Démarrage rapide
Le Notebook suivant montre comment exécuter WordCount à l'aide de Structured Streaming avec Kinesis.
Kinesis WordCount avec Notebook Structured Streaming
Configurer les options Kinesis
Dans Databricks Runtime 13.3 LTS et versions ultérieures, vous pouvez utiliser Trigger.AvailableNow avec Kinesis. Voir Ingestion de Kinesis records en tant que batch incrémentiel.
Dans Databricks Runtime 16.1 et versions ulté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.
Veuillez ne pas basculer entre streamName et streamARN pour une query de streaming active. Databricks ne prend pas en charge la modification de 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 en streaming avec un nouveau répertoire de point de contrôle.
Pour la liste complète des options, voir Kinesis.
Monitoring et alertes à faible latence
Les cas d'usage d'alerte nécessitent une faible latence. Pour minimiser la latence :
- Veuillez vérifier que votre requête de streaming est le seul consommateur du Stream Kinesis afin d'optimiser les performances d'extraction et d'éviter les limites de débit Kinesis.
- Définissez l'option
maxFetchDurationà une petite valeur, par exemple 200 ms, pour traiter les données récupérées le plus rapidement 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 utilisezTrigger.AvailableNow, une petite valeur pourrait entraîner un décalage de votre query par rapport aux enregistrements les plus récents dans le Stream Kinesis. - Définissez l'option
minFetchPeriodsur 210 ms pour récupérer le plus fréquemment possible. - Définissez l’option
shardsPerTaskou configurez le cluster de manière à ce que# cores in cluster >= 2 * (# Kinesis shards) / shardsPerTask. Cela garantit que les tâches de prélecture en arrière-plan et les tâches de query streaming s’exécutent simultanément.
Si votre query reçoit des données toutes les 5 secondes, vous pourriez dépasser les limites de débit Kinesis. Vérifiez vos configurations.
Surveiller les métriques Kinesis
Kinesis signale le nombre de millisecondes qu'un consommateur est en retard 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 Workspaces dans le processus de query en streaming. Consultez le 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 du tableau de bord de progression de la query de streaming. Voici un exemple :
{
"sources": [
{
"description": "KinesisV2[stream]",
"metrics": {
"avgMsBehindLatest": "32000.0",
"maxMsBehindLatest": "32000",
"minMsBehindLatest": "32000"
}
}
]
}
Ingérer des enregistrements Kinesis par batchs incrémentiels
Dans Databricks Runtime 13.3 LTS et versions ultérieures, Databricks prend en charge l'utilisation de Trigger.AvailableNow avec les sources de données Kinesis pour la sémantique de batch incrémentiel. Ce qui suit décrit la configuration de base :
- Lorsqu'une lecture en micro-batch se déclenche en mode « disponible maintenant », l'heure actuelle est enregistrée par le client Databricks.
- Databricks interroge le système source pour tous les enregistrements avec des timestamps entre cette heure enregistrée et le point de contrôle précédent.
- Databricks charge ces enregistrements en utilisant la sémantique
Trigger.AvailableNow.
Databricks utilise un mécanisme de meilleur effort pour tenter de consommer tous les enregistrements qui existent dans les flux Kinesis lorsque la query de streaming s'exécute. En raison de petites différences potentielles dans les Timestamp et d'un manque de garantie dans l'ordre des sources de données, un Trigger pourrait ne pas inclure certains enregistrements. Les enregistrements omis sont traités dans le prochain micro-batch déclenché.
Si la requête ne parvient toujours pas à récupérer d'enregistrements du Kinesis Stream, même s'il y a des enregistrements, essayez d'augmenter la valeur maxFetchDuration.
Consultez AvailableNow: Traitement par batch incrémentiel.
Gérer la perte de données
Utilisez failOnDataLoss uniquement si votre charge de travail 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 l'absence d'enregistrements, redémarrez le Stream avec un nouveau point de contrôle pour retraiter tous les enregistrements.
Databricks vous recommande de l'utiliser uniquement comme solution temporaire à un problème de perte de données. Recherchez et corrigez la cause première, telle qu'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 Kinesis Stream 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 de streaming échouent lorsqu'elles détectent une perte de données potentielle. Pour configurer la query afin d'ignorer les enregistrements illisibles et de continuer le traitement, définissez spark.databricks.kinesis.failOnDataLoss sur false dans votre configuration Spark :
spark.conf.set("spark.databricks.kinesis.failOnDataLoss", "false")
Écrire dans Kinesis
Utilisez l'extrait de code suivant comme ForeachSink pour écrire des données dans Kinesis. Cela nécessite un Dataset[(String, Array[Byte])].
L'extrait de code suivant offre une sémantique au moins une fois , et non exactement une fois.
Notebook de Kinesis Foreach Sink
Recommandations pour réduire la latence avec Kinesis
Cette section contient des recommandations pour le dépannage des diverses causes de latence pour les flux Kinesis.
La source Kinesis exécute des Jobs Spark dans un thread d'arrière-plan pour pré-récupérer périodiquement les données Kinesis, puis met en cache les données dans la mémoire de l'exécuteur Spark. Une fois chaque étape de pré-récupération terminée, la query de streaming peut traiter les données mises en cache. L’étape de prélecture affecte considérablement la latence et le throughput de bout en bout observés.
Réduire la latence de prélecture
Pour optimiser la latence minimale des query et l'utilisation maximale des ressources, utilisez le calcul suivant :
total number of CPU cores in the cluster (across all executors) >= total number of Kinesis shards / shardsPerTask.
minFetchPeriod peut créer plusieurs appels d’API GetRecords vers le shard Kinesis jusqu’à ce qu’il atteigne ReadProvisionedThroughputExceeded. Si une exception se produit, ce n'est peut-être pas un problème car le connecteur maximise l'utilisation du shard Kinesis.
Éviter 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 à partir de Kinesis chaque fois qu'il rencontre une erreur de limitation de débit et enregistre cet événement dans les Logs avec un message : "Hit rate limit. Sleeping for 5 seconds."
Vous pourriez voir ces erreurs pendant qu'un Stream est en cours de rattrapage. Si vous voyez ces erreurs après qu'un Stream est à jour, vous devrez peut-être ajuster la charge de travail, soit en augmentant la capacité Kinesis dans AWS, soit en ajustant les options de prélecture dans Spark.
Évitez le disk spill
Si vous constatez une augmentation soudaine du volume de données dans vos flux Kinesis, la capacité tampon attribuée pourrait se remplir et ne pas se vider assez rapidement pour ajouter de nouvelles données. Spark spill data from the buffer to disk, ce qui ralentit le traitement des Stream, et un événement apparaît dans le journal avec un message comme le suivant :
./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 davantage de nœuds ou en augmentant la mémoire par nœud, ou réduisez le paramètre de configuration fetchBufferSize.
Tâches d'écriture S3 suspendues
Activez la spéculation Spark pour mettre fin aux tâches suspendues qui empêcheraient le traitement en Stream de se poursuivre. Pour s'assurer que les tâches ne soient pas terminées de manière trop agressive, ajustez le quantile et le multiplicateur pour ce paramètre avec soin. Databricks vous recommande de définir spark.speculation.multiplier sur 3 et spark.speculation.quantile sur 0.95 et de les ajuster au besoin.
Réduire la latence due à la création de points de contrôle dans les flux avec état.
Databricks recommande d'utiliser RocksDB avec un point de contrôle du journal des modifications pour les requêtes de streaming avec état. Voir Activer le pointage de contrôle du journal des modifications.
Configurez l’éventail étendu (EFO) de Kinesis pour les lectures de query en streaming
Dans Databricks Runtime 11.3 et versions ultérieures, le connecteur Databricks Runtime Kinesis prend en charge l'utilisation de la fonctionnalité de répartition améliorée (EFO) d'Amazon Kinesis.
Le fan-out amélioré de Kinesis offre 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 au lieu du mode pull.
Par default, une query Structured Streaming configurée en mode EFO s'enregistre en tant que consommateur avec un throughput dédié et un nom de consommateur ainsi qu'un ARN (Amazon Resource Name) de consommateur uniques dans Kinesis Data Stream.
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 remplacer ce comportement. Le consumerName doit être une chaîne qui contient uniquement des lettres, des chiffres et les caractères spéciaux _ . -.
Au redémarrage de la requête, la source Kinesis utilise le mode d'interrogation pour rejouer le dernier lot non validé s'il en existe un. Une fois que le Stream a rejoué le batch non validé, la source revient en mode EFO pour les lectures ultérieures.
Un consommateur EFO enregistré entraîne des frais supplémentaires sur Amazon Kinesis. Pour désinscrire le consommateur automatiquement à l'arrêt de la query, définissez l'option requireConsumerDeregistration sur true. Databricks ne peut pas garantir le désenregistrement lors d'événements tels que les pannes de Driver ou les défaillances de nœuds. En cas d'échec de 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 programme 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 est uniquement disponible en Scala avec le compute défini sur le mode d'accès dédié. Voir les modes d'accès.
Pour utiliser cet utilitaire dans un Notebook Databricks :
-
Dans un nouveau Notebook Databricks attaché à un cluster actif, créez un
AWSKinesisConsumerManageravec les informations d’authentification requises.Scalaimport com.databricks.sql.kinesis.AWSKinesisConsumerManager
val manager = AWSKinesisConsumerManager.newManager()
.option("serviceCredential", serviceCredentialName)
.option("region", kinesisRegion)
.create() -
Répertorier et afficher les consommateurs.
Scalaval consumers = manager.listConsumers("<stream name>")
display(consumers) -
Enregistrer un consommateur pour un Stream donné.
Scalaval consumerARN = manager.registerConsumer("<stream name>", "<consumer name>") -
Désenregistrer un consommateur pour le stream donné.
Scalamanager.deregisterConsumer("<stream name>", "<consumer name>")