Aller au contenu principal

Utiliser Azure Event Hubs comme source de données de pipeline

Vous pouvez traiter les messages d'Azure Event Hubs dans un pipeline à l'aide de l'Endpoint compatible Kafka. Vous ne pouvez pas utiliser le connecteur Structured Streaming Event Hubs car cette bibliothèque ne fait pas partie de Databricks Runtime, et les Lakeflow pipelines ne vous permettent pas d'utiliser des bibliothèques JVM tierces.

Comment un pipeline peut-il se connecter à Azure Event Hubs ?

Azure Event Hubs fournit un endpoint compatible avec Apache Kafka que vous pouvez utiliser avec le connecteur Structured Streaming Kafka, disponible dans Databricks Runtime, pour traiter les messages d'Azure Event Hubs. Pour plus d'informations sur la compatibilité entre Azure Event Hubs et Apache Kafka, consultez Utiliser Azure Event Hubs à partir d'applications Apache Kafka.

Les étapes suivantes décrivent comment connecter un pipeline à une instance Event Hubs existante et consommer des événements à partir d'un sujet. Pour effectuer ces étapes, vous avez besoin des valeurs de connexion Event Hubs suivantes :

  • Le nom de l'espace de noms Event Hubs.
  • Le nom de l'instance Event Hub dans l'espace de noms Event Hubs.
  • Un nom de politique d'accès partagé et une clé de politique pour Event Hubs. Par défaut, une politique RootManageSharedAccessKey est créée pour chaque espace de noms Event Hubs. Cette politique a les autorisations manage, send et listen. Si votre pipeline lit uniquement à partir d'Event Hubs, Databricks recommande de créer une nouvelle politique avec une autorisation d'écoute seulement.

Pour plus d'informations sur la chaîne de connexion Event Hubs, consultez Obtenir une chaîne de connexion Event Hubs.

remarque
  • Azure Event Hubs propose des options OAuth 2.0 et de signature d'accès partagé (SAS) pour autoriser l'accès à vos Ressources sécurisées. Ces instructions utilisent l'authentification basée sur SAS.
  • Si vous obtenez la chaîne de connexion Event Hubs depuis le portail Azure, elle peut ne pas contenir la valeur EntityPath. La valeur EntityPath est requise uniquement lors de l'utilisation du connecteur Event Hubs Structured Streaming. L’utilisation du connecteur Kafka Structured Streaming nécessite de fournir uniquement le nom de la rubrique.

Stockez la clé de la politique dans un secret Databricks

Étant donné que la clé de stratégie est une information sensible, Databricks vous recommande de ne pas coder en dur la valeur dans le code de votre pipeline. Utilisez plutôt les secrets Databricks pour stocker et gérer l'accès à la clé.

L'exemple suivant utilise l'interface CLI de Databricks pour créer un Secret Scope et stocker la clé dans ce Secret Scope. Dans votre code de pipeline, utilisez la fonction dbutils.secrets.get() avec les scope-name et shared-policy-name pour récupérer la valeur de la clé.

Bash
databricks --profile <profile-name> secrets create-scope <scope-name>

databricks --profile <profile-name> secrets put-secret <scope-name> <shared-policy-name> --string-value <shared-policy-key>

Pour plus d'informations sur les secrets Databricks, consultez Gestion des secrets.

Créez un pipeline et ajoutez du code pour consommer des événements

L'exemple suivant lit les événements IoT à partir d'un sujet, mais vous pouvez adapter l'exemple aux exigences de votre application. À titre de bonne pratique, Databricks recommande d'utiliser les paramètres du pipeline pour configurer les variables d'application. Votre code de pipeline utilise ensuite la fonction spark.conf.get() pour récupérer les valeurs.

Python
from pyspark import pipelines as dp
import pyspark.sql.types as T
from pyspark.sql.functions import *

# Event Hubs configuration
EH_NAMESPACE = spark.conf.get("iot.ingestion.eh.namespace")
EH_NAME = spark.conf.get("iot.ingestion.eh.name")

EH_CONN_SHARED_ACCESS_KEY_NAME = spark.conf.get("iot.ingestion.eh.accessKeyName")
SECRET_SCOPE = spark.conf.get("io.ingestion.eh.secretsScopeName")
EH_CONN_SHARED_ACCESS_KEY_VALUE = dbutils.secrets.get(scope = SECRET_SCOPE, key = EH_CONN_SHARED_ACCESS_KEY_NAME)

EH_CONN_STR = f"Endpoint=sb://{EH_NAMESPACE}.servicebus.windows.net/;SharedAccessKeyName={EH_CONN_SHARED_ACCESS_KEY_NAME};SharedAccessKey={EH_CONN_SHARED_ACCESS_KEY_VALUE}"
# Kafka Consumer configuration

KAFKA_OPTIONS = {
"kafka.bootstrap.servers" : f"{EH_NAMESPACE}.servicebus.windows.net:9093",
"subscribe" : EH_NAME,
"kafka.sasl.mechanism" : "PLAIN",
"kafka.security.protocol" : "SASL_SSL",
"kafka.sasl.jaas.config" : f"kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModule required username=\"$ConnectionString\" password=\"{EH_CONN_STR}\";",
"kafka.request.timeout.ms" : spark.conf.get("iot.ingestion.kafka.requestTimeout"),
"kafka.session.timeout.ms" : spark.conf.get("iot.ingestion.kafka.sessionTimeout"),
"maxOffsetsPerTrigger" : spark.conf.get("iot.ingestion.spark.maxOffsetsPerTrigger"),
"failOnDataLoss" : spark.conf.get("iot.ingestion.spark.failOnDataLoss"),
"startingOffsets" : spark.conf.get("iot.ingestion.spark.startingOffsets")
}

# PAYLOAD SCHEMA
payload_ddl = """battery_level BIGINT, c02_level BIGINT, cca2 STRING, cca3 STRING, cn STRING, device_id BIGINT, device_name STRING, humidity BIGINT, ip STRING, latitude DOUBLE, lcd STRING, longitude DOUBLE, scale STRING, temp BIGINT, timestamp BIGINT"""
payload_schema = T._parse_datatype_string(payload_ddl)

# Basic record parsing and adding ETL audit columns
def parse(df):
return (df
.withColumn("records", col("value").cast("string"))
.withColumn("parsed_records", from_json(col("records"), payload_schema))
.withColumn("iot_event_timestamp", expr("cast(from_unixtime(parsed_records.timestamp / 1000) as timestamp)"))
.withColumn("eh_enqueued_timestamp", expr("timestamp"))
.withColumn("eh_enqueued_date", expr("to_date(timestamp)"))
.withColumn("etl_processed_timestamp", col("current_timestamp"))
.withColumn("etl_rec_uuid", expr("uuid()"))
.drop("records", "value", "key")
)

@dp.create_table(
comment="Raw IOT Events",
table_properties={
&quot;quality&quot;: &quot;bronze&quot;,
&quot;pipelines.reset.allowed&quot;: &quot;false&quot; # preserves the data in the delta table if you do full refresh
},
partition_cols=["eh_enqueued_date"]
)
@dp.expect("valid_topic", "topic IS NOT NULL")
@dp.expect("valid records", "parsed_records IS NOT NULL")
def iot_raw():
return (
spark.readStream
.format("kafka")
.options(**KAFKA_OPTIONS)
.load()
.transform(parse)
)

Créer le pipeline

Créez un nouveau pipeline avec un fichier source Python, et entrez le code ci-dessus.

Le code référence les paramètres configurés. Utilisez la configuration JSON suivante, en remplaçant les valeurs d'espaces réservés par les valeurs appropriées pour votre environnement (voir la liste, à la suite du JSON). Vous pouvez définir les parameters en utilisant l’interface utilisateur des paramètres, ou en modifiant directement le JSON des paramètres. Pour plus d'information sur l’utilisation des paramètres de pipeline pour paramétrer votre pipeline, consultez Utiliser des parameters avec les pipelines.

Ce fichier de paramètres définit également l’emplacement de stockage pour un compte de stockage Azure Data Lake Storage (ADLS). En guise de bonne pratique, ce pipeline n'utilise pas le chemin de stockage DBFS par default, mais plutôt un compte de stockage ADLS. Pour plus d'informations sur la configuration de l'authentification d'un compte de stockage ADLS, consultez Accéder en toute sécurité aux informations d'identification de stockage avec des secrets dans un pipeline.

JSON
{
"storage": "abfss://<container-name>@<storage-account-name>.dfs.core.windows.net/iot/",
"configuration": {
"iot.ingestion.eh.namespace": "<eh-namespace>",
"iot.ingestion.eh.accessKeyName": "<eh-policy-name>",
"iot.ingestion.eh.name": "<eventhub>",
"io.ingestion.eh.secretsScopeName": "<secret-scope-name>",
"iot.ingestion.spark.maxOffsetsPerTrigger": "50000",
"iot.ingestion.spark.startingOffsets": "latest",
"iot.ingestion.spark.failOnDataLoss": "false",
"iot.ingestion.kafka.requestTimeout": "60000",
"iot.ingestion.kafka.sessionTimeout": "30000"
}
}

Remplacez les espaces réservés suivants :

  • <container-name> avec le nom d'un conteneur de compte de stockage Azure.
  • <storage-account-name> avec le nom d'un compte de stockage ADLS.
  • <eh-namespace> avec le nom de votre espace de noms Event Hubs.
  • <eh-policy-name> avec la clé du Secret Scope pour la clé de politique Event Hubs.
  • <eventhub> avec le nom de votre instance Event Hubs.
  • <secret-scope-name> avec le nom du Databricks Secret Scope qui contient la clé de politique Event Hubs.