Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
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:
topictopicstopicsPattern
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
valueoszlopba. - 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 egyvalueoszlopba 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()