Skip to main content

risingwave_connector/source/google_pubsub/
mod.rs

1// Copyright 2022 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::HashMap;
16
17use anyhow::Context;
18use google_cloud_gax::conn::Environment;
19use google_cloud_pubsub::apiv1;
20use google_cloud_pubsub::client::google_cloud_auth::credentials::CredentialsFile;
21use google_cloud_pubsub::client::google_cloud_auth::project;
22use google_cloud_pubsub::client::google_cloud_auth::token::DefaultTokenSourceProvider;
23use google_cloud_pubsub::client::{Client, ClientConfig};
24use google_cloud_pubsub::subscriber::SubscriberConfig;
25use google_cloud_pubsub::subscription::Subscription;
26use risingwave_common::bail;
27use risingwave_common::util::env_var::env_var_is_true;
28use serde::Deserialize;
29
30pub mod enumerator;
31pub mod source;
32pub mod split;
33
34pub use enumerator::*;
35use phf::{Set, phf_set};
36use serde_with::{DisplayFromStr, serde_as};
37pub use source::*;
38pub use split::*;
39use with_options::WithOptions;
40
41use crate::connector_common::{DISABLE_DEFAULT_CREDENTIAL, resolve_pubsub_project_id};
42use crate::enforce_secret::EnforceSecret;
43use crate::error::ConnectorResult;
44use crate::source::SourceProperties;
45
46pub const GOOGLE_PUBSUB_CONNECTOR: &str = "google_pubsub";
47
48const DEFAULT_ACK_DEADLINE_SECONDS: i32 = 60;
49// Pub/Sub messages are acknowledged only after a checkpoint. The upstream client default of 50
50// can therefore stall each reader between checkpoints and severely limit throughput.
51const DEFAULT_MAX_OUTSTANDING_MESSAGES: i64 = 1024;
52const DEFAULT_MAX_OUTSTANDING_BYTES: i64 = 1_000_000_000;
53
54/// # Implementation Notes
55/// Pub/Sub does not rely on persisted state (`SplitImpl`) to start from a position.
56/// It rely on Pub/Sub to load-balance messages between all Readers.
57/// We `ack` received messages after checkpoint (see `WaitCheckpointWorker`) to achieve at-least-once delivery.
58#[serde_as]
59#[derive(Clone, Debug, Deserialize, WithOptions)]
60pub struct PubsubProperties {
61    /// The Google Pub/Sub project ID. If omitted, the connector uses the project ID from
62    /// the credentials or Application Default Credentials.
63    #[serde(rename = "pubsub.project_id")]
64    pub project_id: Option<String>,
65
66    /// Pub/Sub subscription to consume messages from.
67    ///
68    /// Note that we rely on Pub/Sub to load-balance messages between all Readers pulling from
69    /// the same subscription. So one `subscription` (i.e., one `Source`) can only used for one MV
70    /// (shared between the actors of its fragment).
71    /// Otherwise, different MVs on the same Source will both receive part of the messages.
72    /// TODO: check and enforce this on Meta.
73    #[serde(rename = "pubsub.subscription")]
74    pub subscription: String,
75
76    /// use the connector with a pubsub emulator
77    /// <https://cloud.google.com/pubsub/docs/emulator>
78    #[serde(rename = "pubsub.emulator_host")]
79    pub emulator_host: Option<String>,
80
81    /// `credentials` is a JSON string containing the service account credentials. If omitted,
82    /// the connector uses Google Application Default Credentials (ADC) when allowed by the
83    /// deployment environment.
84    /// See the [service-account credentials guide](https://developers.google.com/workspace/guides/create-credentials#create_credentials_for_a_service_account).
85    /// The service account must have the `pubsub.subscriber` [role](https://cloud.google.com/pubsub/docs/access-control#roles).
86    #[serde(rename = "pubsub.credentials")]
87    pub credentials: Option<String>,
88
89    /// `start_offset` is a numeric timestamp, ideally the publish timestamp of a message
90    /// in the subscription. If present, the connector will attempt to seek the subscription
91    /// to the timestamp and start consuming from there. Note that the seek operation is
92    /// subject to limitations around the message retention policy of the subscription. See
93    /// [Seeking to a timestamp](https://cloud.google.com/pubsub/docs/replay-overview#seeking_to_a_timestamp) for
94    /// more details.
95    #[serde(rename = "pubsub.start_offset.nanos")]
96    pub start_offset: Option<String>,
97
98    /// `start_snapshot` is a named pub/sub snapshot. If present, the connector will first seek
99    /// to the snapshot before starting consumption. Snapshots are the preferred seeking mechanism
100    /// in pub/sub because they guarantee retention of:
101    /// - All unacknowledged messages at the time of their creation.
102    /// - All messages created after their creation.
103    /// Besides retention guarantees, snapshots are also more precise than timestamp-based seeks.
104    /// See [Seeking to a snapshot](https://cloud.google.com/pubsub/docs/replay-overview#seeking_to_a_timestamp) for
105    /// more details.
106    #[serde(rename = "pubsub.start_snapshot")]
107    pub start_snapshot: Option<String>,
108
109    /// Deprecated: ignored since adaptive split mode was introduced.
110    /// Split count now adapts automatically to the number of actors.
111    /// Kept for backward compatibility with existing DDL.
112    #[serde_as(as = "Option<DisplayFromStr>")]
113    #[serde(rename = "pubsub.parallelism")]
114    pub parallelism: Option<u32>,
115
116    /// The ack deadline in seconds for the streaming pull subscriber.
117    /// This is the maximum time the server will wait for an ack before redelivering the message.
118    /// Must be between 10 and 600 seconds. Defaults to 60.
119    #[serde_as(as = "Option<DisplayFromStr>")]
120    #[serde(rename = "pubsub.ack_deadline_seconds")]
121    #[with_option(allow_alter_on_fly)]
122    pub ack_deadline_seconds: Option<i32>,
123
124    /// The maximum number of unacknowledged messages delivered to each streaming pull reader.
125    /// Pub/Sub pauses delivery to a reader when this limit is reached. Must be greater than 0.
126    /// Defaults to 1024.
127    #[serde_as(as = "Option<DisplayFromStr>")]
128    #[serde(rename = "pubsub.max_outstanding_messages")]
129    #[with_option(allow_alter_on_fly)]
130    pub max_outstanding_messages: Option<i64>,
131
132    /// The maximum total size of unacknowledged messages delivered to each streaming pull reader.
133    /// Pub/Sub pauses delivery to a reader when this limit is reached. Must be greater than 0.
134    /// Defaults to 1 GB.
135    #[serde_as(as = "Option<DisplayFromStr>")]
136    #[serde(rename = "pubsub.max_outstanding_bytes")]
137    #[with_option(allow_alter_on_fly)]
138    pub max_outstanding_bytes: Option<i64>,
139
140    #[serde(flatten)]
141    pub unknown_fields: HashMap<String, String>,
142}
143
144impl EnforceSecret for PubsubProperties {
145    const ENFORCE_SECRET_PROPERTIES: Set<&'static str> = phf_set! {
146        "pubsub.credentials",
147    };
148}
149
150impl SourceProperties for PubsubProperties {
151    type Split = PubsubSplit;
152    type SplitEnumerator = PubsubSplitEnumerator;
153    type SplitReader = PubsubSplitReader;
154
155    const SOURCE_NAME: &'static str = GOOGLE_PUBSUB_CONNECTOR;
156}
157
158impl crate::source::UnknownFields for PubsubProperties {
159    fn unknown_fields(&self) -> HashMap<String, String> {
160        self.unknown_fields.clone()
161    }
162}
163
164impl PubsubProperties {
165    pub(crate) fn subscriber_config(&self) -> ConnectorResult<SubscriberConfig> {
166        let stream_ack_deadline_seconds = self
167            .ack_deadline_seconds
168            .unwrap_or(DEFAULT_ACK_DEADLINE_SECONDS);
169        if !(10..=600).contains(&stream_ack_deadline_seconds) {
170            bail!("pubsub.ack_deadline_seconds must be between 10 and 600");
171        }
172
173        let max_outstanding_messages = self
174            .max_outstanding_messages
175            .unwrap_or(DEFAULT_MAX_OUTSTANDING_MESSAGES);
176        if max_outstanding_messages <= 0 {
177            bail!("pubsub.max_outstanding_messages must be greater than 0");
178        }
179
180        let max_outstanding_bytes = self
181            .max_outstanding_bytes
182            .unwrap_or(DEFAULT_MAX_OUTSTANDING_BYTES);
183        if max_outstanding_bytes <= 0 {
184            bail!("pubsub.max_outstanding_bytes must be greater than 0");
185        }
186
187        Ok(SubscriberConfig {
188            stream_ack_deadline_seconds,
189            max_outstanding_messages,
190            max_outstanding_bytes,
191            ..Default::default()
192        })
193    }
194
195    pub(crate) async fn subscription_client(&self) -> ConnectorResult<Subscription> {
196        let auth_config = project::Config::default()
197            .with_audience(apiv1::conn_pool::AUDIENCE)
198            .with_scopes(&apiv1::conn_pool::SCOPES);
199        let (environment, detected_project_id) = if let Some(credentials) = &self.credentials {
200            let credentials = CredentialsFile::new_from_str(credentials)
201                .await
202                .context("failed to parse Google Cloud Pub/Sub credentials")?;
203            let provider = DefaultTokenSourceProvider::new_with_credentials(
204                auth_config,
205                Box::new(credentials),
206            )
207            .await
208            .context("failed to initialize Google Cloud Pub/Sub token source")?;
209            let project_id = provider.project_id.clone();
210            (Environment::GoogleCloud(Box::new(provider)), project_id)
211        } else if let Some(emulator_host) = &self.emulator_host {
212            (Environment::Emulator(emulator_host.clone()), None)
213        } else {
214            if env_var_is_true(DISABLE_DEFAULT_CREDENTIAL) {
215                bail!(
216                    "Google Application Default Credentials are disabled; configure `pubsub.credentials` or `pubsub.emulator_host`"
217                );
218            }
219
220            let provider = DefaultTokenSourceProvider::new(auth_config)
221                .await
222                .context(
223                    "failed to initialize Google Cloud Pub/Sub ADC; provide `pubsub.credentials`, configure ADC, or use `pubsub.emulator_host`",
224                )?;
225            let project_id = provider.project_id.clone();
226            (Environment::GoogleCloud(Box::new(provider)), project_id)
227        };
228
229        let project_id = resolve_pubsub_project_id(
230            self.project_id.as_deref(),
231            detected_project_id.as_deref(),
232            matches!(&environment, Environment::Emulator(_)),
233        )
234        .context(
235            "Google Cloud Pub/Sub project ID is unavailable; configure `pubsub.project_id` or provide credentials/ADC with a project ID",
236        )?;
237        let config = ClientConfig {
238            environment,
239            project_id: Some(project_id),
240            ..Default::default()
241        };
242        let client = Client::new(config)
243            .await
244            .context("error initializing pubsub client")?;
245
246        Ok(client.subscription(&self.subscription))
247    }
248}
249
250#[cfg(test)]
251mod tests {
252    use serde_json::json;
253
254    use super::*;
255
256    fn parse_pubsub_properties(extra: serde_json::Value) -> PubsubProperties {
257        let mut value = json!({
258            "pubsub.subscription": "projects/test/subscriptions/test",
259            "pubsub.emulator_host": "localhost:8900",
260        });
261        value
262            .as_object_mut()
263            .unwrap()
264            .extend(extra.as_object().unwrap().clone());
265        serde_json::from_value(value).unwrap()
266    }
267
268    #[test]
269    fn test_subscriber_config_defaults() {
270        let config = parse_pubsub_properties(json!({}))
271            .subscriber_config()
272            .unwrap();
273
274        assert_eq!(config.stream_ack_deadline_seconds, 60);
275        assert_eq!(config.max_outstanding_messages, 1024);
276        assert_eq!(config.max_outstanding_bytes, 1_000_000_000);
277    }
278
279    #[test]
280    fn test_subscriber_config_overrides() {
281        let config = parse_pubsub_properties(json!({
282            "pubsub.ack_deadline_seconds": "120",
283            "pubsub.max_outstanding_messages": "2048",
284            "pubsub.max_outstanding_bytes": "1048576",
285        }))
286        .subscriber_config()
287        .unwrap();
288
289        assert_eq!(config.stream_ack_deadline_seconds, 120);
290        assert_eq!(config.max_outstanding_messages, 2048);
291        assert_eq!(config.max_outstanding_bytes, 1_048_576);
292    }
293
294    #[test]
295    fn test_subscriber_config_validation() {
296        let invalid_values = [
297            (
298                json!({"pubsub.ack_deadline_seconds": "9"}),
299                "pubsub.ack_deadline_seconds must be between 10 and 600",
300            ),
301            (
302                json!({"pubsub.max_outstanding_messages": "0"}),
303                "pubsub.max_outstanding_messages must be greater than 0",
304            ),
305            (
306                json!({"pubsub.max_outstanding_bytes": "0"}),
307                "pubsub.max_outstanding_bytes must be greater than 0",
308            ),
309        ];
310
311        for (value, expected_error) in invalid_values {
312            let error = parse_pubsub_properties(value)
313                .subscriber_config()
314                .unwrap_err();
315            assert!(error.to_string().contains(expected_error));
316        }
317    }
318}