Stream depuis Apache Pulsar
Aperçu
Cette fonctionnalité est en aperçu public.
Avec Databricks Runtime 14.1 et versions supérieures, vous pouvez utiliser Structured Streaming pour Stream des données depuis Apache Pulsar sur Databricks.
Structured Streaming fournit une sémantique de traitement exactement une fois pour les données lues à partir des sources Pulsar.
Exemple de syntaxe
Voici un exemple de base d'utilisation du Structured Streaming pour lire depuis Pulsar :
- Python
- Scala
query = (spark.readStream
.format("pulsar")
.option("service.url", "pulsar://broker.example.com:6650")
.option("topics", "topic1,topic2")
.load()
)
val query = spark.readStream
.format("pulsar")
.option("service.url", "pulsar://broker.example.com:6650")
.option("topics", "topic1,topic2")
.load()
Pour lire à partir de rubriques Pulsar, vous devez fournir un service.url et l'une des options suivantes :
topictopicstopicsPattern
Pour obtenir une liste complète des options, consultez Configurer les options de lecture en streaming de Pulsar.
S'authentifier auprès de Pulsar
Databricks prend en charge l'authentification truststore et keystore pour Pulsar. Databricks recommande que vous utilisiez des secrets pour stocker les détails de configuration.
Pour la liste complète des options d'authentification, consultez Authentification.
Exemple
L'exemple suivant illustre la configuration des options d'authentification :
- Python
- Scala
client_auth_params = dbutils.secrets.get(scope="pulsar", key="clientAuthParams")
client_pw = dbutils.secrets.get(scope="pulsar", key="clientPw")
# clientAuthParams is a comma-separated list of key-value pairs, such as:
# "keyStoreType:JKS,keyStorePath:/var/private/tls/client.keystore.jks,keyStorePassword:clientpw"
query = (spark.readStream
.format("pulsar")
.option("service.url", "pulsar://broker.example.com:6650")
.option("topics", "topic1,topic2")
.option("startingOffsets", starting_offsets)
.option("pulsar.client.authPluginClassName", "org.apache.pulsar.client.impl.auth.AuthenticationKeyStoreTls")
.option("pulsar.client.authParams", client_auth_params)
.option("pulsar.client.useKeyStoreTls", "true")
.option("pulsar.client.tlsTrustStoreType", "JKS")
.option("pulsar.client.tlsTrustStorePath", trust_store_path)
.option("pulsar.client.tlsTrustStorePassword", client_pw)
.load()
)
val clientAuthParams = dbutils.secrets.get(scope = "pulsar", key = "clientAuthParams")
val clientPw = dbutils.secrets.get(scope = "pulsar", key = "clientPw")
// clientAuthParams is a comma-separated list of key-value pairs, such as:
// "keyStoreType:JKS,keyStorePath:/var/private/tls/client.keystore.jks,keyStorePassword:clientpw"
val query = spark.readStream
.format("pulsar")
.option("service.url", "pulsar://broker.example.com:6650")
.option("topics", "topic1,topic2")
.option("startingOffsets", startingOffsets)
.option("pulsar.client.authPluginClassName", "org.apache.pulsar.client.impl.auth.AuthenticationKeyStoreTls")
.option("pulsar.client.authParams", clientAuthParams)
.option("pulsar.client.useKeyStoreTls", "true")
.option("pulsar.client.tlsTrustStoreType", "JKS")
.option("pulsar.client.tlsTrustStorePath", trustStorePath)
.option("pulsar.client.tlsTrustStorePassword", clientPw)
.load()
schéma Pulsar
Lorsque vous lisez à partir de Pulsar, le schéma des lignes dépend des schémas des topics source.
- Pour les rubriques avec un schéma Avro ou JSON, les noms et types de champs sont conservés dans le DataFrame Spark résultant.
- Pour les rubriques sans schéma ou avec un type de données simple dans Pulsar, la charge utile est chargée dans une colonne
value. - Si vous configurez le Stream pour lire plusieurs rubriques avec des schémas différents, définissez
allowDifferentTopicSchemaspour charger le contenu brut dans une colonnevalue.
Les enregistrements Pulsar ont les champs de métadonnées suivants :
Colonne | Type |
|---|---|
|
|
|
|
|
|
|
|
|
|
|
|
Configurer les options pour la lecture en streaming Pulsar
Pour la liste complète des options, consultez Pulsar.
Construction de JSON d'offsets de départ
Pour utiliser un ID de message personnalisé qui spécifie un décalage, au format JSON, avec l'option startingOffsets, consultez l'exemple suivant :
import org.apache.spark.sql.pulsar.JsonUtils
import org.apache.pulsar.client.api.MessageId
import org.apache.pulsar.client.impl.MessageIdImpl
val topic = "my-topic"
val msgId: MessageId = new MessageIdImpl(ledgerId, entryId, partitionIndex)
val startOffsets = JsonUtils.topicOffsets(Map(topic -> msgId))
query = spark.readStream
.format("pulsar")
.option("service.url", "pulsar://broker.example.com:6650")
.option("topic", topic)
.option("startingOffsets", startOffsets)
.load()