Adatfolyam az Apache Pulsar rendszerből

Fontos

Ez a funkció a nyilvános előzetes verzióban érhető el.

A Databricks Runtime 14.1-ben és újabb verziókban a strukturált streamelés használatával streamelheti az adatokat az Apache Pulsarból az Azure Databricksen.

A strukturált streamelés pontosan egyszeri feldolgozási szemantikát biztosít a Pulsar-forrásokból beolvasott adatokhoz.

Szintaxispélda

Az alábbiakban egy egyszerű példa következik a Structured Streaming használatára, Pulsarból történő beolvasáshoz:

Python

query = (spark.readStream
  .format("pulsar")
  .option("service.url", "pulsar://broker.example.com:6650")
  .option("topics", "topic1,topic2")
  .load()
)

Scala

val query = spark.readStream
  .format("pulsar")
  .option("service.url", "pulsar://broker.example.com:6650")
  .option("topics", "topic1,topic2")
  .load()

A Pulsar-témákból való olvasáshoz meg kell adnia egy service.url elemet és a következő lehetőségek egyikét:

  • topic
  • topics
  • topicsPattern

A beállítások teljes listáját a Pulsar-streamelési olvasás beállításainak konfigurálása című témakörben találja.

Hitelesítés Pulsarban

Az Azure Databricks támogatja a truststore és a keystore hitelesítést a Pulsar számára. A Databricks azt javasolja, hogy titkos kulcsokkal tárolja a konfiguráció részleteit.

A hitelesítési lehetőségek teljes listáját a Hitelesítés című témakörben találja.

Example

Az alábbi példa a hitelesítési beállítások konfigurálását mutatja be:

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

Scala

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

Pulsar-séma

Amikor a Pulsarból olvas, a sorok sémája a forrás témaköreinek sémáitól függ.

  • Az Avro- vagy JSON-sémát tartalmazó témakörök esetében a mezőnevek és a mezőtípusok megmaradnak az eredményül kapott Spark DataFrame-ben.
  • Séma nélküli vagy egyszerű adattípusú témakörök esetén a fizikai adattartalom betöltésre kerül egy value oszlopba.
  • Ha úgy konfigurálja a streamet, hogy több témát olvasson különböző sémákkal, állítsa allowDifferentTopicSchemas értékre, hogy a nyers tartalom egy value oszlopba töltődjön be.

A Pulsar-rekordok a következő metaadatmezőket tartalmaznak:

oszlop Típus
__key binary
__topic string
__messageId binary
__publishTime timestamp
__eventTime timestamp
__messageProperties map<String, String>

A Pulsar streaminges olvasási beállítások konfigurálása

A lehetőségek teljes listáját a Pulsarban találja.

JSON kezdeti eltolások létrehozása

Ha egy eltolást meghatározó egyéni üzenetazonosítót szeretne használni JSON-ként a startingOffsets beállítással, tekintse meg a következő példát:

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