Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Используйте встроенный соединитель для подписки на Google Pub/Sub. Этот коннектор поддерживает семантику однократной обработки строк, полученных от подписчика.
Примечание.
Pub/Sub может публиковать повторяющиеся записи, либо записи могут поступать подписчику не по порядку. Необходимо написать код для обработки повторяющихся и устаревших строк.
Настройка потока Pub/Sub
В следующем примере кода показано, как настроить структурированную потоковую передачу из pub/Sub и пройти проверку подлинности с помощью закрытых ключей.
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
Метрики хода выполнения Structured Streaming отражают количество полученных и готовых к обработке строк, размер полученных и готовых к обработке строк, а также количество дубликатов, обнаруженных с момента запуска потока.
Ниже приведен пример метрик Pub/Sub:
"metrics" : {
"numDuplicatesSinceStreamStart" : "1",
"numRecordsReadyToProcess" : "1",
"sizeOfRecordsReadyToProcess" : "8"
}
Ограничения
Pub/Sub не поддерживает спекулятивное выполнение с spark.speculation.