使用内置连接器订阅 Google Pub/Sub。 此连接器对来自订阅者的行具有“仅处理一次”的处理语义。
注意
Pub/Sub 可能会发布重复的行数据,或者这些行数据到达订阅方时可能是乱序的。 必须编写代码来处理重复行和无序行。
配置 Pub/Sub 流
以下代码示例演示了如何配置从 Pub/Sub 读取 Structured Streaming 数据,并使用私钥进行身份验证。
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')
);
有关更多配置选项,请参阅“配置 Pub/Sub 流式处理读取”选项。
配置对 Pub/Sub 的访问权限
您的凭据必须具备以下角色:
| 角色 | 必需或可选 | 角色的使用方式 |
|---|---|---|
roles/pubsub.viewer 或 roles/viewer |
必填 | 检查订阅是否存在,然后获取订阅。 |
roles/pubsub.subscriber |
必填 | 从订阅中提取数据。 |
roles/pubsub.editor 或 roles/editor |
可选 | 如果订阅不存在,则启用创建订阅的功能,并在流终止时使用deleteSubscriptionOnStreamStop删除订阅。 |
注意
如果你是在资源级别而不是项目级别授予 roles/pubsub.viewer 和 roles/pubsub.subscriber,则必须将这两个角色同时应用于主题和订阅。 如果您未使用可选的 roles/pubsub.editor 或 roles/editor 角色,仅在主题上授予必需的角色是不够的。
Databricks 建议在使用密钥时使用机密。 授权连接需要以下选项:
clientEmailclientIdprivateKeyprivateKeyId
了解 Pub/Sub 架构
流的模式与从 Pub/Sub 获取的行相匹配,如下表所述:
| 字段 | 类型 |
|---|---|
messageId |
StringType |
payload |
ArrayType[ByteType] |
attributes |
StringType |
publishTimestampInMillis |
LongType |
配置 Pub/Sub 流式读取选项
某些 Pub/Sub 配置选项使用提取概念,而不是微批处理。 这是一个内部实现细节,这些选项的工作方式与其他结构化流处理连接器类似,只是行会先被提取出来,然后再进行处理。
有关选项的完整列表,请参阅 Pub/Sub。
将增量批处理与 Pub/Sub 配合使用
您可以使用 Trigger.AvailableNow 以增量批处理的方式从 Pub/Sub 源中消费可用的行。
Azure Databricks 在 Trigger.AvailableNow 设置中记录你开始读取的时间戳。 该批次处理的行包括所有先前获取的数据,以及时间戳小于记录的开始时间戳的任何新发布的行。 有关详细信息,请参阅 AvailableNow:增量批处理。
监控 Pub/Sub 流式指标
结构化流进度指标报告已提取且可供处理的行数、这些已提取且可供处理的行的大小,以及自流开始以来发现的重复项数量。
下面是 Pub/Sub 指标的示例:
"metrics" : {
"numDuplicatesSinceStreamStart" : "1",
"numRecordsReadyToProcess" : "1",
"sizeOfRecordsReadyToProcess" : "8"
}
限制
Pub/Sub 不支持使用 spark.speculation 进行推测执行。