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
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
- Scala
- SQL
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()
)
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()
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')
);
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 d’utiliser des secrets lorsque vous utilisez des clés. Les options suivantes sont requises pour autoriser une connexion :
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.