FAQ
Questions fréquemment posées sur l'utilisation de Kafka avec Databricks.
Pourquoi est-ce que j'obtiens une erreur indiquant qu'une option Kafka n'est pas prise en charge ou n'est pas reconnue ?
Cette erreur se produit si vous oubliez d'utiliser le préfixe kafka. lors de la configuration des options client Kafka. Toutes les options transmises directement au client Kafka doivent être préfixées par kafka.:
Le code suivant présente des options incorrectes auxquelles il manque le préfixe kafka. :
.option("security.protocol", "SASL_SSL")
.option("sasl.mechanism", "PLAIN")
Le code suivant présente les options correctes :
.option("kafka.security.protocol", "SASL_SSL")
.option("kafka.sasl.mechanism", "PLAIN")
Les options du connecteur Spark Kafka (comme subscribe, startingOffsets, maxOffsetsPerTrigger) ne nécessitent pas de préfixe. Pour la liste complète des options, consultez Kafka.
Pourquoi est-ce que j'obtiens une erreur concernant les classes Kafka occultées ?
Databricks exige l'utilisation de classes Kafka masquées (préfixées par kafkashaded. ou shadedmskiam.). Si vous voyez des erreurs, telles que RESTRICTED_STREAMING_OPTION_PERMISSION_ENFORCED, vous devez utiliser les noms de classes masquées :
org.apache.kafka.*les classes nécessitent le préfixekafkashaded.. Par exemple :kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModulesoftware.amazon.msk.*les classes nécessitent le préfixeshadedmskiam.. Par exemple :shadedmskiam.software.amazon.msk.auth.iam.IAMLoginModule
Pourquoi obtenez-vous un TimeoutException lors de la connexion à Kafka ?
Parmi les causes courantes, citons :
- Connectivité réseau : Le cluster de compute ne peut pas atteindre les brokers Kafka. Vérifiez les règles de pare-feu, les groupes de sécurité et les configurations de Virtual Private Cloud (VPC).
- Serveurs d'amorçage incorrects : Vérifiez que le Hostname et le port
kafka.bootstrap.serverssont corrects. - Résolution DNS : vérifiez que les hostnames du broker Kafka peuvent être résolus depuis le réseau Databricks.
- Problèmes SSL/TLS : si vous utilisez SSL, vérifiez que les certificats sont correctement configurés.
Pour les configurations Private Link ou d'appairage VPC, vérifiez que les routes réseau correctes sont en place.
Dois-je utiliser le mode batch ou streaming pour Kafka ?
Cela dépend de votre cas d'utilisation :
- Mode streaming (
spark.readStream) : Utilisez-le lorsque vous avez besoin d'un traitement de données continu ou d'une ingestion à faible latence. - Mode batch
spark.read() : à utiliser pour les chargements de données uniques, les remplissages ou le debugging. Nécessite à la foisstartingOffsetsetendingOffsets.
Consultez Configurer les intervalles Trigger Structured Streaming pour plus de détails sur la configuration des intervalles Trigger tels que AvailableNow, ProcessingTime et le mode temps réel.
Puis-je lire à partir de plusieurs rubriques Kafka dans un seul Stream ?
Oui, vous pouvez utiliser :
subscribe: Fournissez une liste de sujets séparés par des virgules, par.option("subscribe", "topic1,topic2")exemple.subscribePattern: Utilisez un modèle d'expression régulière Java pour faire correspondre les noms de rubrique, par exemple.option("subscribePattern", "topic-.*").
Comment utiliser Kafka avec les LakeFlow Pipelines ?
Les LakeFlow Pipelines disposent d’un support intégré pour les sources Kafka.
Vous pouvez définir une table de streaming qui lit à partir de Kafka, comme dans le code suivant :
- Python
- SQL
import dlt
@dlt.table
def kafka_bronze():
return (spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:port>")
.option("subscribe", "<topic>")
.load()
)
CREATE OR REFRESH STREAMING TABLE kafka_bronze AS
SELECT * FROM STREAM read_kafka(
bootstrapServers => '<server:port>',
subscribe => '<topic>'
);
Consultez Charger des données dans les pipelines pour plus de détails sur les sources de streaming dans les LakeFlow Pipelines.
Comment désérialiser les colonnes de clé et de valeur Kafka ?
Les colonnes key et value sont renvoyées en tant que type BINARY. Utilisez les opérations DataFrame pour les désérialiser en fonction de votre format de données :
- **Données de type chaîne** : utilisez
cast("string")pour convertir le binaire en chaîne. - Données JSON : utilisez
from_json()après conversion en chaîne de caractères. Voirfrom_jsonfonction. - Données Avro : Utilisez
from_avro()pour désérialiser les données encodées en Avro. Voir Lire et écrire des données Avro en streaming. - **Buffers de protocole** : Utilisez
from_protobuf()pour désérialiser les données protobuf. Consultez Lire et écrire des buffers de protocole.
Pourquoi est-ce que j'obtiens une erreur d'écriture idempotente ?
Databricks Runtime 13.3 LTS et versions ultérieures inclut une version plus récente de la bibliothèque kafka-clients qui permet les écritures idempotentes by default. Si votre cluster Kafka utilise la version 2.8.0 ou inférieure avec des ACL configurées mais sans IDEMPOTENT_WRITE activé, l'écriture échoue avec : org.apache.kafka.common.KafkaException: Cannot execute transactional method because we are in an error state.
Résolvez cette erreur en effectuant une mise à niveau vers la version 2.8.0 ou ultérieure de Kafka, ou en définissant .option("kafka.enable.idempotence", "false") lors de la configuration de votre writer Structured Streaming.
Qu'est-ce que KAFKA_DATA_LOSS_ERROR et comment puis-je le résoudre ?
Cette erreur se produit lorsque la source Kafka détecte que les décalages (offsets) stockés dans le point de contrôle ne sont plus disponibles dans Kafka, généralement parce que :
- Le stream a été mis en pause plus longtemps que la période de rétention Kafka.
- Les données du sujet Kafka ont été supprimées ou le sujet a été recréé.
- Le broker Kafka a subi une perte de données.
Pour résoudre :
- Si la perte de données est acceptable : définissez
.option("failOnDataLoss", "false")pour permettre au Stream de continuer à partir du décalage disponible le plus ancien. - Si la perte de données est inacceptable : Reset le point de contrôle et retraitez à partir des décalages
earliest, ou restaurez les données Kafka manquantes.
Pour plus d'informations, consultez KAFKA_DATA_LOSS error condition.
Comment contrôler le débit auquel les données sont lues depuis Kafka ?
Utilisez l'option maxOffsetsPerTrigger pour limiter le nombre de décalages (environ le nombre d'enregistrements) traités par micro-batch. Cela permet d'éviter les grands batchs qui pourraient submerger le traitement en aval ou causer des problèmes de mémoire lors du rattrapage d'un arriéré.
- Python
- Scala
- SQL
df = (spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:port>")
.option("subscribe", "<topic>")
.option("maxOffsetsPerTrigger", 10000)
.load()
)
val df = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:port>")
.option("subscribe", "<topic>")
.option("maxOffsetsPerTrigger", 10000)
.load()
SELECT * FROM STREAM read_kafka(
bootstrapServers => '<server:port>',
subscribe => '<topic>',
maxOffsetsPerTrigger => '10000'
);
Vous pouvez également utiliser des options telles que minPartitions ou maxRecordsPerPartition pour contrôler le nombre de partitions Spark créées pour chaque batch.
Comment puis-je surveiller le retard de mon stream par rapport aux derniers offsets Kafka ?
Utilisez les métriques avgOffsetsBehindLatest, maxOffsetsBehindLatest et minOffsetsBehindLatest disponibles dans la progression de la query de streaming. Ceux-ci indiquent combien de décalages votre Stream a de retard par rapport au dernier décalage disponible sur toutes les partitions de sujet souscrites. Consultez le monitoring des requêtes Structured Streaming sur Databricks.
Vous pouvez également utiliser estimatedTotalBytesBehindLatest pour estimer le nombre total d'octets de données qui n'ont pas encore été traités.
Pourquoi mes métriques de décalage d'offset Kafka affichent-elles des valeurs non nulles persistantes après la mise à niveau vers Databricks Runtime 17.1 ?
Dans Databricks Runtime 17.1 et versions ultérieures, les derniers offsets Kafka sont récupérés après l'achèvement de chaque micro-batch. Sur les rubriques qui reçoivent continuellement des données, les métriques de backlog peuvent afficher des valeurs non nulles, faibles et persistantes. Il s'agit d'un comportement attendu et n'indique pas que le Stream prend du retard.
Dans Databricks Runtime 17.0 et versions antérieures, les derniers offsets Kafka sont récupérés à l'heure de start du micro-batch. Les métriques de backlog peuvent renvoyer 0 lorsque les queries en streaming consomment systématiquement tous les enregistrements disponibles au start du micro-batch.
Si les valeurs sont grandes ou augmentent continuellement, le stream pourrait ne pas suivre les données entrantes. Consultez monitoring Structured Streaming queries sur Databricks.
Pourquoi l'initialisation de mon Stream Kafka est-elle lente ?
Les flux Kafka nécessitent du temps pour :
- Connectez-vous au cluster Kafka et récupérez les métadonnées.
- Découvrir les partitions de sujet.
- Récupérer les décalages initiaux.
Pour les clusters Kafka on-premise ou distants, la latence du réseau peut avoir un impact significatif sur le temps d'initialisation. Si vous exécutez des pipelines Trigger/planifiés avec des redémarrages fréquents, envisagez d'utiliser le mode de streaming continu pour éviter la charge de travail d'initialisation répétée.
Pourquoi l'ajout de plus d'exécuteurs Spark n'augmente-t-il pas mon throughput Kafka ?
Une fois les brokers Kafka saturés, l'ajout de davantage d'exécuteurs Spark augmente les coûts sans augmenter le throughput.
Signes que Kafka est le goulot d'étranglement :
- Le throughput stagne malgré l'ajout de cœurs.
- L'utilisation du CPU ou du réseau du broker Kafka est élevée.
- Les tâches Spark se terminent rapidement, mais attendent de nouvelles données.
Pour résoudre ce problème, mettez à l'échelle votre cluster Kafka en ajoutant des brokers ou en augmentant le nombre de partitions pour répartir la charge.
Comment puis-je optimiser les coûts et l'utilisation du compute pour le streaming Kafka ?
Pour les modes micro-batch et AvailableNow :
- Dimensionnez correctement votre cluster : surveillez les métriques et définissez une taille de cluster fixe appropriée pour la charge maximale.
- Utiliser
maxOffsetsPerTrigger: Limitez les tailles de batch pour contrôler l'utilisation des Ressources lors des pics de charge. - **Évitez l'autoscaling** : les jobs de streaming s'exécutent en continu, et l'ajout ou la suppression de nœuds entraîne des frais généraux de rééquilibrage des tâches.
- Réduire l'asymétrie des données : les partitions asymétriques entraînent un traitement significativement plus important de données pour certaines tâches que pour d'autres, ce qui crée des retardataires qui ralentissent l'achèvement global du batch et gaspillent les ressources de compute pour des tâches inactives. Utilisez l'option
minPartitionspour diviser les grandes partitions Kafka en plus petites partitions Spark pour un traitement plus équilibré.
Pour le mode en temps réel, le dimensionnement du compute est particulièrement important car les tâches peuvent rester inactives en attendant les données. Points clés à considérer :
- Définissez
maxPartitionsafin que chaque tâche gère plusieurs partitions Kafka pour réduire la surcharge. - Optimisez
spark.sql.shuffle.partitionspour les Jobs gourmands en shuffle.
Consultez Dimensionnement de compute pour obtenir des conseils sur le dimensionnement des clusters pour le mode temps réel.
Pourquoi mon Stream ne renvoie-t-il aucun enregistrement alors que des données existent dans le sujet ?
Parmi les causes courantes, citons :
- Paramètre
startingOffsetsincorrect : la valeur default estlatest, qui ne lit que les nouvelles données arrivant après le start du Stream. DéfinissezstartingOffsetssurearliestpour lire les données existantes. - Nom du sujet incorrect : Vérifiez que vous êtes abonné au bon sujet.
- **Problèmes d'authentification** : Votre Stream a peut-être été connecté avec succès mais ne dispose pas des autorisations nécessaires pour lire le sujet. Vérifiez vos ACL Kafka.
- Expiration des offsets : Si votre Stream a été arrêté pendant une longue période et que les offsets du checkpoint ont expiré (ont été supprimés par la rétention Kafka), vous pourriez avoir besoin de reset le checkpoint ou d'ajuster
failOnDataLoss.