Aller au contenu principal

Se connecter à Apache Kafka

Cette page décrit comment vous pouvez utiliser Apache Kafka comme source ou comme récepteur lors de l’exécution de charges de travail Structured Streaming sur Databricks.

Pour plus d'informations sur Kafka, consultez la documentation Apache Kafka.

Lire les données depuis Kafka

Utilisez le format kafka pour configurer les connexions à Kafka. Voici un exemple de lecture en streaming :

Python
df = (spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("subscribe", "<topic>")
.option("startingOffsets", "latest")
.load()
)

Databricks prend également en charge les lectures par batch depuis Kafka, comme dans l'exemple suivant :

Python
df = (spark.read
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("subscribe", "<topic>")
.option("startingOffsets", "earliest")
.option("endingOffsets", "latest")
.load()
)

Pour le chargement incrémental par batch, Databricks recommande d'utiliser Kafka avec Trigger.AvailableNow. Voir AvailableNow: Traitement par batch incrémental.

Dans Databricks Runtime 13.3 LTS et versions ultérieures, Databricks fournit également une fonction SQL pour la lecture des données Kafka. Le streaming avec SQL n'est pris en charge que dans les Lakeflow Pipelines ou avec des tables de streaming dans Databricks SQL. Consultez la fonction à valeur tabulaireread_kafka.

Configurer le lecteur Structured Streaming de Kafka

Pour les queries batch et streaming, vous devez définir les serveurs d'amorçage pour la source Kafka avec l'option suivante :

Clé

Valeur

Description

kafka.bootstrap.servers

Une liste d'hôtes séparés par des virgules

Les serveurs de bootstrap du cluster Kafka

Clé

Valeur

Description

kafka.bootstrap.servers

Une liste d'hôtes séparés par des virgules

Les serveurs de bootstrap du cluster Kafka

Pour définir les sujets d'abonnement, vous devez spécifier l'une des options suivantes :

Option

Valeur

Description

subscribe

Une liste de rubriques séparées par des virgules.

La liste des sujets auxquels s'abonner.

subscribePattern

Chaîne regex Java.

Le modèle utilisé pour s'abonner à un ou plusieurs sujets.

assign

Chaîne JSON {"topicA":[0,1],"topic":[2,4]}.

topicPartitions spécifique à consommer.

Option

Valeur

Description

subscribe

Une liste de rubriques séparées par des virgules.

La liste des sujets auxquels s'abonner.

subscribePattern

Chaîne regex Java.

Le modèle utilisé pour s'abonner à un ou plusieurs sujets.

assign

Chaîne JSON {"topicA":[0,1],"topic":[2,4]}.

topicPartitions spécifique à consommer.

Consultez Kafka pour la liste complète des options disponibles.

Schéma pour les lignes Kafka

Le lecteur Kafka Structured Streaming renvoie des lignes avec le schéma suivant :

Colonne

Type

key

binary

value

binary

topic

string

partition

int

offset

long

timestamp

timestamp

timestampType

int

Colonne

Type

key

binary

value

binary

topic

string

partition

int

offset

long

timestamp

timestamp

timestampType

int

Le key et le value sont toujours désérialisés en tant que tableaux d’octets avec le ByteArrayDeserializer. Utilisez les opérations de DataFrame (telles que cast("string") ou from_avro) pour désérialiser explicitement les clés et les valeurs.

Écrire des données vers Kafka

Voici un exemple d'écriture en streaming vers Kafka :

Python
(df.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("topic", "<topic>")
.start()
)

Databricks prend également en charge les sémantiques d'écriture par batch vers les récepteurs de données Kafka, comme le montre l'exemple suivant :

Python
(df.write
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("topic", "<topic>")
.save()
)

Configurer l’enregistreur Kafka Structured Streaming

important

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 un récepteur Kafka utilise la version 2.8.0 ou antérieure avec des ACL configurées, mais sans l'activation de IDEMPOTENT_WRITE, l'écriture échoue avec le message d'erreur 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.

Voici les options courantes pour les écritures vers Kafka :

Clé

Valeur

Valeur par défaut

Description

kafka.boostrap.servers

Une liste séparée par des virgules de <host:port>

Aucun

Obligatoire. La configuration Kafka bootstrap.servers.

topic

STRING

non défini

Facultatif. Définit la rubrique pour toutes les lignes à écrire. Cette option remplace toute colonne de rubrique existante dans les données.

includeHeaders

BOOLEAN

false

Facultatif. S'il faut inclure les en-têtes Kafka dans la ligne.

Clé

Valeur

Valeur par défaut

Description

kafka.boostrap.servers

Une liste séparée par des virgules de <host:port>

Aucun

Obligatoire. La configuration Kafka bootstrap.servers.

topic

STRING

non défini

Facultatif. Définit la rubrique pour toutes les lignes à écrire. Cette option remplace toute colonne de rubrique existante dans les données.

includeHeaders

BOOLEAN

false

Facultatif. S'il faut inclure les en-têtes Kafka dans la ligne.

Voir Kafka sink pour la liste complète des options disponibles.

Schéma pour l'enregistreur Kafka

Lors de l'écriture de données vers Kafka, le DataFrame fourni peut inclure les champs suivants :

Nom de colonne

Obligatoire ou facultatif

Type

key

Facultatif

STRING OU BINARY

value

Obligatoire

STRING OU BINARY

headers

Facultatif

ARRAY

topic

facultatif (ignoré si topic est défini comme option d’écriture)

STRING

partition

Facultatif

INT

Nom de colonne

Obligatoire ou facultatif

Type

key

Facultatif

STRING OU BINARY

value

Obligatoire

STRING OU BINARY

headers

Facultatif

ARRAY

topic

facultatif (ignoré si topic est défini comme option d’écriture)

STRING

partition

Facultatif

INT

Authentification

Databricks prend en charge plusieurs méthodes d'authentification pour Kafka, y compris les identifiants de service Unity Catalog, SASL/SSL, et les options spécifiques au cloud pour AWS MSK, Azure Event Hub et Google Cloud Managed Kafka. Voir Authentification.

Récupérez les métriques Kafka

Pour surveiller le décalage derrière Kafka pour une query streaming, utilisez les métriques avgOffsetsBehindLatest, maxOffsetsBehindLatest et minOffsetsBehindLatest. Ces métriques signalent le décalage moyen, maximal et minimal pour toutes les partitions de sujets souscrites, par rapport aux décalages les plus récents dans Kafka. Consultez Lecture interactive des métriques.

remarque

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 au moment du `start` du micro-`batch`. Les métriques de backlog peuvent retourner 0 lorsque les query en streaming consomment de manière constante tous les enregistrements disponibles au start du micro-batch.

Pour estimer les données restantes à lire pour une query, utilisez la métrique estimatedTotalBytesBehindLatest. Cette métrique estime le nombre total d'octets restants sur toutes les partitions abonnées en fonction des batchs traités au cours des 300 dernières secondes. Vous pouvez modifier la fenêtre de temps utilisée pour cette estimation en définissant l'option bytesEstimateWindowLength.

Par exemple, pour définir la longueur de la fenêtre à 10 minutes :

Python
df = (spark.readStream
.format("kafka")
.option("bytesEstimateWindowLength", "10m") # m for minutes, you can also use "600s" for 600 seconds
)

Si vous exécutez le Stream dans un Notebook, vous pouvez voir ces métriques sous l’onglet **Données brutes** du tableau de bord de progression du streaming query :

JSON
{
"sources": [
{
"description": "KafkaV2[Subscribe[topic]]",
"metrics": {
"avgOffsetsBehindLatest": "4.0",
"maxOffsetsBehindLatest": "4",
"minOffsetsBehindLatest": "4",
"estimatedTotalBytesBehindLatest": "80.0"
}
}
]
}

Consultez Monitoring des queries Structured Streaming sur Databricks pour plus d'information.

Exemple de Kafka à Delta Lake

L'exemple suivant montre un workflow complet pour une écriture en streaming incrémentielle de Kafka vers une table Delta Lake à l'aide du availableNow trigger. Vous pouvez utiliser cette approche pour les charges de travail d'ingestion de données incrémentielles.

Cet exemple utilise un schéma JSON fixe. Pour les autres formats comme Avro ou Protobuf, utilisez from_avro ou from_protobuf. Vous pouvez également l'intégrer à un registre de schémas. Voir Exemple avec le registre de schémas.

Python
from pyspark.sql.functions import from_json, col

# Define simple JSON schemas for key and value
key_schema = "user_id STRING"
value_schema = "event_type STRING, event_ts TIMESTAMP"

# Configure Kafka options with service credentials
kafka_options = {
"kafka.bootstrap.servers": "<bootstrap-server>:9092",
"subscribe": "<topic-name>",
"databricks.serviceCredential": "<service-credential-name>",
}

# Read from Kafka and parse JSON
parsed_df = (spark.readStream
.format("kafka")
.options(**kafka_options)
.load()
.select(
from_json(col("key").cast("string"), key_schema).alias("key"),
from_json(col("value").cast("string"), value_schema).alias("value")
)
.select("key.*", "value.*")
)

# Write to Delta table
query = (parsed_df.writeStream
.format("delta")
.option("checkpointLocation", "/path/to/checkpoint")
.trigger(availableNow=True)
.toTable("catalog.schema.events_table")
)

query.awaitTermination()
remarque

Sur le compute serverless Databricks, le trigger availableNow est recommandé pour le streaming incrémentiel. Pour le streaming continu à faible latence, utilisez le mode continu des Lakeflow pipelines. Consultez les triggers de streaming structuré pour la liste complète des options prises en charge.