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é.
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
- Scala
- SQL
query = (spark.readStream
.format("pubsub")
.option("subscriptionId", "mysub")
.option("topicId", "mytopic")
.option("projectId", "myproject")
.option("serviceCredential", "service-credential-name")
.load()
)
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")
.option("serviceCredential", "service-credential-name")
.load()
CREATE OR REFRESH STREAMING TABLE pubsub_raw
AS SELECT * FROM STREAM read_pubsub(
subscriptionId => 'mysub',
projectId => 'myproject',
topicId => 'mytopic',
serviceCredential => 'service-credential-name'
);
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 |
|---|---|---|
| Obligatoire | Vérifie si l'abonnement existe et récupère l'abonnement. |
| Obligatoire | Récupère les données d'un abonnement. |
| Facultatif | Permet la création d'un abonnement s'il n'existe pas et permet l'utilisation du |
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.
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 :
clientEmailclientIdprivateKeyprivateKeyId
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 |
|---|---|
|
|
|
|
|
|
|
|
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 :
"metrics" : {
"numDuplicatesSinceStreamStart" : "1",
"numRecordsReadyToProcess" : "1",
"sizeOfRecordsReadyToProcess" : "8"
}
Limitations
Pub/Sub ne prend pas en charge l'exécution spéculative avec spark.speculation.