Se connecter à Apache Kafka

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.