Aller au contenu principal

Stream depuis Apache Pulsar

info

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
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 :

  • topic
  • topics
  • topicsPattern

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
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()
)

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 allowDifferentTopicSchemas pour charger le contenu brut dans une colonne value.

Les enregistrements Pulsar ont les champs de métadonnées suivants :

Colonne

Type

__key

binary

__topic

string

__messageId

binary

__publishTime

timestamp

__eventTime

timestamp

__messageProperties

map<String, String>

Colonne

Type

__key

binary

__topic

string

__messageId

binary

__publishTime

timestamp

__eventTime

timestamp

__messageProperties

map<String, String>

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 :

Scala
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()