Aller au contenu principal

S'abonner à Google Pub/Sub

Utilisez le connecteur intégré pour vous abonner à Google Pub/Sub. Ce connecteur possède une sémantique de traitement à exécution unique pour les lignes de l'abonné.

remarque

Pub/Sub pourrait publier des lignes dupliquées, ou les lignes pourraient arriver à l'abonné dans le désordre. Vous devez écrire du code pour gérer les lignes dupliquées et désordonnées.

Configure un Stream Pub/Sub

L'exemple de code suivant montre comment configurer une lecture Structured Streaming à partir de Pub/Sub et s'authentifier avec des clés privées.

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()
)

Pour plus d’options de configuration, consultez Configurer les options de lecture en streaming Pub/Sub.

Configurer l'accès à Pub/Sub

Vos identifiants doivent avoir les rôles suivants :

Rôles

Obligatoire ou facultatif

Utilisation du rôle

roles/pubsub.viewer OU roles/viewer

Obligatoire

Vérifie si l'abonnement existe et récupère l'abonnement.

roles/pubsub.subscriber

Obligatoire

Récupère les données d'un abonnement.

roles/pubsub.editor OU roles/editor

Facultatif

Permet la création d'un abonnement s'il n'existe pas et permet l'utilisation du deleteSubscriptionOnStreamStop pour supprimer les abonnements à la fin du Stream.

Rôles

Obligatoire ou facultatif

Utilisation du rôle

roles/pubsub.viewer OU roles/viewer

Obligatoire

Vérifie si l'abonnement existe et récupère l'abonnement.

roles/pubsub.subscriber

Obligatoire

Récupère les données d'un abonnement.

roles/pubsub.editor OU roles/editor

Facultatif

Permet la création d'un abonnement s'il n'existe pas et permet l'utilisation du deleteSubscriptionOnStreamStop pour supprimer les abonnements à la fin du Stream.

remarque

Si vous accordez roles/pubsub.viewer et roles/pubsub.subscriber au niveau de la ressource plutôt qu'au niveau du projet, vous devez appliquer les deux rôles au sujet et à l'abonnement. Si vous n'utilisez pas les rôles facultatifs roles/pubsub.editor ou roles/editor, l'octroi des rôles requis sur le sujet seul n'est pas suffisant.

Databricks vous recommande d’utiliser des secrets lorsque vous utilisez des clés. Les options suivantes sont requises pour autoriser une connexion :

  • clientEmail
  • clientId
  • privateKey
  • privateKeyId

Comprendre le schéma Pub/Sub

Le schéma du stream correspond aux lignes récupérées depuis Pub/Sub, comme décrit dans le tableau suivant :

Champ

Type

messageId

StringType

payload

ArrayType[ByteType]

attributes

StringType

publishTimestampInMillis

LongType

Champ

Type

messageId

StringType

payload

ArrayType[ByteType]

attributes

StringType

publishTimestampInMillis

LongType

Configurer les options pour la lecture en streaming Pub/Sub

Certaines options de configuration Pub/Sub utilisent le concept de récupérations au lieu de micro-batchs . Ceci est un détail d'implémentation interne, et les options fonctionnent de manière similaire aux autres connecteurs Structured Streaming, sauf que les lignes sont récupérées puis traitées.

Pour la liste complète des options, voir Pub/Sub.

Utilisez le traitement par batch incrémentiel avec Pub/Sub

Vous pouvez utiliser Trigger.AvailableNow pour consommer les lignes disponibles à partir des sources Pub/Sub comme un batch incrémentiel.

Databricks enregistre le Timestamp lorsque vous commencez une lecture avec le paramètre Trigger.AvailableNow. Les lignes traitées par le batch incluent toutes les données précédemment récupérées et toutes les nouvelles lignes publiées avec un Timestamp inférieur au Timestamp de start enregistré. Pour plus d'information, consultez AvailableNow: traitement par batch incrémentiel.

Surveiller les métriques de streaming Pub/Sub

Les métriques de progression de Structured Streaming indiquent le nombre de lignes récupérées et prêtes à être traitées, la taille des lignes récupérées et prêtes à être traitées, et le nombre de doublons vus depuis le start du Stream.

Voici un exemple de métriques Pub/Sub :

JSON
"metrics" : {
"numDuplicatesSinceStreamStart" : "1",
"numRecordsReadyToProcess" : "1",
"sizeOfRecordsReadyToProcess" : "8"
}

Limitations

Pub/Sub ne prend pas en charge l'exécution spéculative avec spark.speculation.