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

Les exemples suivants montrent comment lire depuis Pub/Sub et s'authentifier avec un identifiant de service. Pour toutes les options d'authentification, voir Configurer l'accès à Pub/Sub.

Python
query = (spark.readStream
.format("pubsub")
.option("subscriptionId", "mysub")
.option("topicId", "mytopic")
.option("projectId", "myproject")
.option("serviceCredential", "service-credential-name")
.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 de configurer un identifiant de service pour les lectures Pub/Sub. Les identifiants de service pour Pub/Sub nécessitent Databricks Runtime 16.1 ou une version ultérieure. Voir Créer des informations d'identification de service. Cependant, si les identifiants de service Databricks ne sont pas disponibles, vous pouvez utiliser directement un compte de service Google (GSA). Les autorisations pour le GSA sont disponibles pour toutes les queries exécutées sur ce cluster. Consultez le compte de service Google.

remarque

Vous ne pouvez pas attacher un GSA à un compute configuré avec le mode d'accès standard.

Configurez les options suivantes pour utiliser un GSA avec un Stream :

  • 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.