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
- Scala
- SQL
df = (spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("subscribe", "<topic>")
.option("startingOffsets", "latest")
.load()
)
val df = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("subscribe", "<topic>")
.option("startingOffsets", "latest")
.load()
CREATE OR REFRESH STREAMING TABLE <table_name> AS
SELECT * FROM STREAM read_kafka(
bootstrapServers => '<server:ip>',
subscribe => '<topic>'
);
Databricks prend également en charge les lectures par batch depuis Kafka, comme dans l'exemple suivant :
- Python
- Scala
- SQL
df = (spark.read
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("subscribe", "<topic>")
.option("startingOffsets", "earliest")
.option("endingOffsets", "latest")
.load()
)
val df = spark.read
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("subscribe", "<topic>")
.option("startingOffsets", "earliest")
.option("endingOffsets", "latest")
.load()
SELECT * FROM read_kafka(
bootstrapServers => '<server:ip>',
subscribe => '<topic>',
startingOffsets => 'earliest',
endingOffsets => 'latest'
);
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 |
|---|---|---|
| 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 |
|---|---|---|
| Une liste de rubriques séparées par des virgules. | La liste des sujets auxquels s'abonner. |
| Chaîne regex Java. | Le modèle utilisé pour s'abonner à un ou plusieurs sujets. |
| Chaîne JSON |
|
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 |
|---|---|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
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
- Scala
(df.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("topic", "<topic>")
.start()
)
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
- Scala
(df.write
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("topic", "<topic>")
.save()
)
df.write
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("topic", "<topic>")
.save()
Configurer l’enregistreur Kafka Structured Streaming
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 |
|---|---|---|---|
| Une liste séparée par des virgules de | Aucun | Obligatoire. La configuration Kafka |
|
| 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. |
|
|
| 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 |
|---|---|---|
| Facultatif |
|
| Obligatoire |
|
| Facultatif |
|
| facultatif (ignoré si |
|
| Facultatif |
|
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.
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
- Scala
df = (spark.readStream
.format("kafka")
.option("bytesEstimateWindowLength", "10m") # m for minutes, you can also use "600s" for 600 seconds
)
val 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 :
{
"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
- Scala
- SQL
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()
import org.apache.spark.sql.functions.{from_json, col}
import org.apache.spark.sql.streaming.Trigger
// Define JSON schemas for key and value
val keySchema = "user_id STRING"
val valueSchema = "event_type STRING, event_ts TIMESTAMP"
// Configure Kafka options with service credentials
val kafkaOptions = Map(
"kafka.bootstrap.servers" -> "<bootstrap-server>:9092",
"subscribe" -> "<topic-name>",
"databricks.serviceCredential" -> "<service-credential-name>"
)
// Read from Kafka and parse JSON
val parsedDF = spark.readStream
.format("kafka")
.options(kafkaOptions)
.load()
.select(
from_json(col("key").cast("string"), keySchema).alias("key"),
from_json(col("value").cast("string"), valueSchema).alias("value")
)
.select("key.*", "value.*")
// Write to Delta table
val query = parsedDF.writeStream
.format("delta")
.option("checkpointLocation", "/path/to/checkpoint")
.trigger(Trigger.ProcessingTime("10 seconds"))
.toTable("catalog.schema.events_table")
query.awaitTermination()
-- Create a streaming table from Kafka using read_kafka
CREATE OR REFRESH STREAMING TABLE catalog.schema.events_table AS
SELECT
key::string:user_id AS user_id,
value::string:event_type AS event_type,
to_timestamp(value::string:event_ts) AS event_ts
FROM STREAM read_kafka(
bootstrapServers => '<bootstrap-server>:9092',
subscribe => '<topic-name>',
serviceCredential => '<service-credential-name>'
);
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.