Подписка на Google Pub/Sub

Используйте встроенный соединитель для подписки на 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 рекомендует использовать секреты при использовании ключей. Для авторизации подключения требуются следующие параметры:

  • clientEmail
  • clientId
  • privateKey
  • privateKeyId

Общие сведения о схеме 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.