Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
Google Pub/Sub'a abone olmak için yerleşik bağlayıcıyı kullanın. Bu bağlayıcı, aboneden gelen satırlar için tam olarak bir kez işleme semantiğine sahiptir.
Not
Pub/Sub yinelenen satırlar yayımlayabilir veya satırlar aboneye sırasız ulaşabilir. Yinelenen ve sırası bozulmuş satırları işlemek için kod yazmanız gerekir.
Pub/Sub akışı yapılandırma
Aşağıdaki kod örneği, Pub/Sub’dan Structured Streaming okumasının nasıl yapılandırılacağını ve özel anahtarlarla nasıl kimlik doğrulanacağını gösterir.
Python
auth_options = {
"clientId": client_id,
"clientEmail": client_email,
"privateKey": private_key,
"privateKeyId": private_key_id
}
query = (spark.readStream
.format("pubsub")
.option("subscriptionId", "mysub")
.option("topicId", "mytopic")
.option("projectId", "myproject")
.options(auth_options)
.load()
)
Scala
val authOptions: Map[String, String] =
Map("clientId" -> clientId,
"clientEmail" -> clientEmail,
"privateKey" -> privateKey,
"privateKeyId" -> privateKeyId)
val query = spark.readStream
.format("pubsub")
// Creates a Pub/Sub subscription if one does not already exist with this ID
.option("subscriptionId", "mysub")
.option("topicId", "mytopic")
.option("projectId", "myproject")
.options(authOptions)
.load()
SQL
CREATE OR REFRESH STREAMING TABLE pubsub_raw
AS SELECT * FROM STREAM read_pubsub(
subscriptionId => 'mysub',
projectId => 'myproject',
topicId => 'mytopic',
clientEmail => secret('pubsub-scope', 'clientEmail'),
clientId => secret('pubsub-scope', 'clientId'),
privateKeyId => secret('pubsub-scope', 'privateKeyId'),
privateKey => secret('pubsub-scope', 'privateKey')
);
Daha fazla yapılandırma seçeneği için Pub/Sub akış okuma seçeneklerini yapılandırma bölümüne bakın.
Pub/Sub erişimi yapılandırma
Kimlik bilgilerinizde aşağıdaki roller bulunmalıdır:
| Roller | Gerekli veya isteğe bağlı | Rol nasıl kullanılır? |
|---|---|---|
roles/pubsub.viewer veya roles/viewer |
Zorunlu | Aboneliğin mevcut olup olmadığını denetler ve aboneliği alır. |
roles/pubsub.subscriber |
Zorunlu | Abonelikten veri getirir. |
roles/pubsub.editor veya roles/editor |
İsteğe bağlı | Abonelik yoksa oluşturulmasını ve akış sonlandırmada abonelikleri silmek için deleteSubscriptionOnStreamStop'i kullanmayı etkinleştirir. |
Not
roles/pubsub.viewer ve roles/pubsub.subscriber rollerini proje düzeyi yerine kaynak düzeyinde verirseniz, her iki rolü de hem konuya hem de aboneliğe uygulamanız gerekir. İsteğe bağlı roles/pubsub.editor veya roles/editor rolleri kullanmıyorsanız, yalnızca konu başlığında gerekli rolleri vermek yeterli değildir.
Databricks, anahtarları kullanırken gizli bilgiler kullanmanızı önerir. Bağlantıyı yetkilendirmek için aşağıdaki seçenekler gereklidir:
clientEmailclientIdprivateKeyprivateKeyId
Pub/Sub şemasını anlama
Akışın şeması, aşağıdaki tabloda açıklandığı gibi Pub/Sub’dan alınan satırlarla eşleşir:
| Alan | Tür |
|---|---|
messageId |
StringType |
payload |
ArrayType[ByteType] |
attributes |
StringType |
publishTimestampInMillis |
LongType |
Pub/Sub akış okuma seçeneklerini yapılandırma
Bazı Pub/Sub yapılandırma seçenekleri, getirme kavramını mikro toplu işlemler yerine kullanır. Bu, uygulamanın iç işleyişine ilişkin bir ayrıntıdır ve seçenekler, satırların önce alınıp ardından işlenmesi dışında, diğer Structured Streaming bağlayıcılarına benzer şekilde çalışır.
Seçeneklerin tam listesi için bkz. Pub/Sub.
Pub/Sub ile artımlı toplu işlem kullanma
Pub/Sub kaynaklarındaki kullanılabilir satırları artımlı bir toplu iş olarak işlemek için Trigger.AvailableNow kullanabilirsiniz.
Azure Databricks, Trigger.AvailableNow ayarıyla bir okumaya başladığınızda zaman damgasını kaydeder. Toplu işlem tarafından işlenen satırlar, daha önce getirilen tüm verileri ve kaydedilen başlangıç zaman damgasından daha az zaman damgası olan yeni yayımlanan satırları içerir. Daha fazla bilgi için bkz AvailableNow. Artımlı toplu işlem.
Pub/Sub akış ölçümlerini izleme
Yapılandırılmış Akış ilerleme durumu ölçümleri getirilen ve işlemeye hazır satır sayısını, getirilen ve işlemeye hazır satırların boyutunu ve akış başlangıcından bu yana görülen yinelenenlerin sayısını bildirir.
Aşağıda Pub/Sub ölçümlerine bir örnek verilmiştir:
"metrics" : {
"numDuplicatesSinceStreamStart" : "1",
"numRecordsReadyToProcess" : "1",
"sizeOfRecordsReadyToProcess" : "8"
}
Sınırlamalar
Pub/Sub, spark.speculation ile spekülatif yürütmeyi desteklemez.