Remarque
L’accès à cette page requiert une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page requiert une autorisation. Vous pouvez essayer de modifier des répertoires.
Cette page explique comment utiliser Apache Kafka comme source ou récepteur lors de l’exécution de charges de travail Structured Streaming sur Azure Databricks.
Pour plus d’informations sur Kafka, consultez la documentation Apache Kafka.
Lire les données de Kafka
Utilisez le kafka format pour configurer les connexions à Kafka. Voici un exemple de lecture en continu :
Python
df = (spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("subscribe", "<topic>")
.option("startingOffsets", "latest")
.load()
)
Scala
val df = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("subscribe", "<topic>")
.option("startingOffsets", "latest")
.load()
SQL
CREATE OR REFRESH STREAMING TABLE <table_name> AS
SELECT * FROM STREAM read_kafka(
bootstrapServers => '<server:ip>',
subscribe => '<topic>'
);
Azure Databricks prend également en charge les lectures par lots à partir de 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()
)
Scala
val df = spark.read
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("subscribe", "<topic>")
.option("startingOffsets", "earliest")
.option("endingOffsets", "latest")
.load()
SQL
SELECT * FROM read_kafka(
bootstrapServers => '<server:ip>',
subscribe => '<topic>',
startingOffsets => 'earliest',
endingOffsets => 'latest'
);
Pour le chargement par lots incrémentiel, Databricks recommande d’utiliser Kafka avec Trigger.AvailableNow. Voir AvailableNow: Traitement par lots incrémentiel.
Dans Databricks Runtime 13.3 LTS et versions ultérieures, Azure Databricks fournit également une fonction SQL pour lire les données Kafka. La diffusion en continu avec SQL est prise en charge uniquement dans les pipelines Lakeflow ou avec des tables de streaming dans Databricks SQL. Consultez read_kafkaTVF.
Configurer le lecteur Kafka de flux structuré
Pour les requêtes de traitement par lots et de diffusion en continu, vous devez définir les serveurs bootstrap pour la source Kafka avec l’option suivante :
| Clé | Valeur | Description |
|---|---|---|
kafka.bootstrap.servers |
Liste séparée par des virgules de host :port | Serveurs de démarrage de cluster Kafka |
Pour définir des rubriques d’abonnement, vous devez spécifier l’une des options suivantes :
| Choix | Valeur | Description |
|---|---|---|
subscribe |
Liste séparée par des virgules des rubriques. | Liste de rubriques auxquelles s’abonner. |
subscribePattern |
Java chaîne regex. | Modèle utilisé pour s’abonner à une ou plusieurs rubriques. |
assign |
Chaîne JSON {"topicA":[0,1],"topic":[2,4]}. |
Spécifique topicPartitions à consommer. |
Consultez Kafka pour obtenir la liste complète des options disponibles.
Schéma pour les lignes Kafka
Le lecteur Kafka Structured Streaming retourne des lignes avec le schéma suivant :
| Colonne | Type |
|---|---|
key |
binary |
value |
binary |
topic |
string |
partition |
int |
offset |
long |
timestamp |
timestamp |
timestampType |
int |
Les key et value sont toujours désérialisées en tant que tableaux d’octets avec le ByteArrayDeserializer. Utilisez des opérations dataFrame (telles que cast("string") ou from_avro) pour désérialiser explicitement les clés et les valeurs.
Écrire des données dans Kafka
Voici un exemple ci-après pour une écriture en streaming dans Kafka :
Python
(df.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("topic", "<topic>")
.start()
)
Scala
df.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("topic", "<topic>")
.start()
Azure Databricks prend également en charge la sémantique d’écriture par lots dans les récepteurs de données Kafka, comme illustré dans l’exemple suivant :
Python
(df.write
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("topic", "<topic>")
.save()
)
Scala
df.write
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("topic", "<topic>")
.save()
Configurer la fonction d’écriture Kafka de flux structuré
Important
Databricks Runtime 13.3 LTS et versions ultérieures incluent une version plus récente de la bibliothèque kafka-clients qui active par défaut les écritures idempotentes. Si un récepteur Kafka utilise la version 2.8.0 ou antérieure avec des listes de contrôle d’accès configurées, mais sans 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.
Pour résoudre cette erreur, mettez à jour vers Kafka version 2.8.0 ou supérieure, ou définissez l’option .option(“kafka.enable.idempotence”, “false”) pendant la configuration de la fonction d’écriture de flux structurés.
Voici les options courantes pour les écritures vers Kafka :
| Clé | Valeur | Valeur par défaut | Description |
|---|---|---|---|
kafka.boostrap.servers |
Une liste des machines virtuelles séparée par des virgules de <host:port> |
aucune | Obligatoire. La configuration Kafka bootstrap.servers. |
topic |
STRING |
non défini | Optional. Définit le sujet de toutes les lignes à écrire. Cette option remplace toute colonne de rubrique qui existe dans les données. |
includeHeaders |
BOOLEAN |
false |
Optional. Indique s’il faut inclure les en-têtes Kafka dans la ligne. |
Consultez Kafka sink pour la liste complète des options disponibles.
Schéma pour le rédacteur Kafka
Lors de l’écriture de données dans Kafka, le DataFrame fourni peut inclure les champs suivants :
| Nom de colonne | Obligatoire ou facultatif | Type |
|---|---|---|
key |
optionnel |
STRING ou BINARY |
value |
obligatoire |
STRING ou BINARY |
headers |
optionnel | ARRAY |
topic |
facultatif (ignoré si topic est défini comme option d’enregistreur) |
STRING |
partition |
optionnel | INT |
Authentification
Azure Databricks prend en charge plusieurs méthodes d’authentification pour Kafka, notamment les informations d’identification du service Unity Catalog, SASL/SSL et des options spécifiques au cloud pour AWS MSK, Azure Event Hubs et Google Cloud Managed Kafka. Consultez Authentification.
Récupérer des métriques Kafka
Pour surveiller le décalage par rapport à Kafka pour une requête de streaming, utilisez les métriques avgOffsetsBehindLatest, maxOffsetsBehindLatest et minOffsetsBehindLatest. Ces métriques indiquent le décalage moyen, maximal et minimal sur l’ensemble des partitions des rubriques souscrites, par rapport aux décalages les plus récents dans Kafka. Consultez Lecture des métriques de manière interactive.
Note
Dans Databricks Runtime 17.1 et versions ultérieures, les derniers offsets Kafka sont récupérés une fois chaque micro-lot terminé. Dans les rubriques qui reçoivent en continu des données, les métriques du backlog peuvent afficher des valeurs petites et persistantes non nulles. Ce comportement est attendu et n’indique pas que le flux est en retard.
Dans Databricks Runtime 17.0 et versions antérieures, les derniers décalages Kafka sont extraits au moment du démarrage du micro-lot. Les métriques du backlog peuvent renvoyer 0 lorsque les requêtes de streaming consomment constamment tous les enregistrements disponibles au début du micro-lot.
Pour estimer les données restantes d’une requête à lire, utilisez la estimatedTotalBytesBehindLatest métrique. Cette métrique estime le nombre total d’octets restants sur toutes les partitions abonnées en fonction des lots 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 durée de la fenêtre sur 10 minutes :
Python
df = (spark.readStream
.format("kafka")
.option("bytesEstimateWindowLength", "10m") # m for minutes, you can also use "600s" for 600 seconds
)
Scala
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 flux dans un notebook, vous pouvez voir ces métriques sous l’onglet Données brutes du tableau de bord de progression des requêtes de diffusion en continu :
{
"sources": [
{
"description": "KafkaV2[Subscribe[topic]]",
"metrics": {
"avgOffsetsBehindLatest": "4.0",
"maxOffsetsBehindLatest": "4",
"minOffsetsBehindLatest": "4",
"estimatedTotalBytesBehindLatest": "80.0"
}
}
]
}
Pour plus d’informations, consultez la surveillance des requêtes Structured Streaming sur Azure Databricks.
Exemple pour Kafka vers Delta Lake
L’exemple suivant montre un flux de travail complet pour une écriture incrémentielle en streaming depuis Kafka vers une table Delta Lake à l’aide du déclencheur availableNow. 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 d’autres formats comme Avro ou Protobuf, utilisez from_avro ou from_protobuf. Vous pouvez également intégrer un registre de schémas. Consultez l’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()
Scala
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()
SQL
-- 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>'
);
Note
Sur le calcul Databricks Serverless, le déclencheur availableNow est recommandé pour le streaming incrémentiel. Pour la diffusion continue à faible latence, utilisez le mode continu des pipelines Lakeflow. Consultez les déclencheurs de streaming structuré pour obtenir la liste complète des options prises en charge.