Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions src/bigquery/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ wkt.workspace = true
gaxi = { workspace = true, features = ["_internal-common", "_internal-grpc-client", "_internal-http-client"] }
tokio = { workspace = true, features = ["time"] }
rust_decimal.workspace = true
uuid.workspace = true

[dev-dependencies]
anyhow.workspace = true
Expand Down
7 changes: 7 additions & 0 deletions src/bigquery/src/query/execution.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ use crate::error::QueryError;
use crate::query::{Query, Result};
use google_cloud_bigquery_v2::client::JobService;
use google_cloud_bigquery_v2::model::{InsertJobRequest, PostQueryRequest};
use google_cloud_gax::options::RequestOptionsBuilder as _;
use std::sync::Arc;

pub(crate) struct PostQueryExecutor {
Expand All @@ -40,6 +41,9 @@ impl PostQueryExecutor {
let res = self
.job_service
.query()
// requests to jobs.query are idempotent because every request
// carries a generated request_id.
.with_idempotency(true)
.with_request(self.request)
.send()
.await?;
Expand Down Expand Up @@ -97,6 +101,9 @@ impl InsertJobExecutor {
.job_service
.insert_job()
.with_request(self.request)
// jobs.insert is idempotent because every request
// carries a generated UUID job_id.
.with_idempotency(true)
.send()
.await?;

Expand Down
78 changes: 69 additions & 9 deletions src/bigquery/src/query/run_query.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,13 @@ use crate::query::{Query, Result};
use google_cloud_bigquery_v2::client::JobService;
use google_cloud_bigquery_v2::model::query_request::JobCreationMode;
use google_cloud_bigquery_v2::model::{
InsertJobRequest, Job, JobConfiguration, PostQueryRequest, QueryRequest,
InsertJobRequest, Job, JobConfiguration, JobReference, PostQueryRequest, QueryRequest,
};
use std::sync::Arc;
use uuid::Uuid;

pub(crate) const JOB_ID_PREFIX: &str = "job_";
pub(crate) const QUERY_REQUEST_ID_PREFIX: &str = "req_";

/// A unified request builder for configuring and running a SQL query.
/// It automatically routes to either `jobs.query` (fast path) or `jobs.insert` (job path)
Expand Down Expand Up @@ -67,8 +71,11 @@ impl RunQuery {

if self.request.force_job_path() {
// Route to jobs.insert
let job_ref = generate_job_reference(&project_id, &self.request.location);
let job_config: JobConfiguration = self.request.into();
let job = Job::new().set_configuration(job_config);
let job = Job::new()
.set_configuration(job_config)
.set_job_reference(job_ref);
let req = InsertJobRequest::new()
.set_job(job)
.set_project_id(project_id);
Expand All @@ -77,12 +84,15 @@ impl RunQuery {
.execute()
.await
} else {
let query_request_id = generate_prefixed_id(QUERY_REQUEST_ID_PREFIX);
// Route to jobs.query
let query_request: QueryRequest = self.request.into();
let query_request = query_request.set_format_options(
google_cloud_bigquery_v2::model::DataFormatOptions::new()
.set_use_int64_timestamp(true),
);
let query_request = query_request
.set_format_options(
google_cloud_bigquery_v2::model::DataFormatOptions::new()
.set_use_int64_timestamp(true),
)
.set_request_id(query_request_id);
let req = PostQueryRequest::new()
.set_project_id(project_id)
.set_query_request(query_request);
Expand All @@ -94,6 +104,23 @@ impl RunQuery {
}
}

fn generate_job_reference(project_id: &str, location: &str) -> JobReference {
let job_id = generate_prefixed_id(JOB_ID_PREFIX);
let mut job_ref = JobReference::new()
.set_project_id(project_id.to_string())
.set_job_id(job_id);

if !location.is_empty() {
job_ref = job_ref.set_location(location.to_string());
}

job_ref
}

fn generate_prefixed_id(prefix: &str) -> String {
format!("{prefix}{}", Uuid::new_v4().simple())
}

include!("../generated/run_query_builder.rs");

#[cfg(test)]
Expand All @@ -105,6 +132,7 @@ mod tests {
Job, JobConfiguration, JobReference, JobStatus, QueryRequest, QueryResponse,
};
use google_cloud_gax::response::Response;
use uuid::Uuid;

type TestResult = anyhow::Result<()>;

Expand Down Expand Up @@ -141,10 +169,38 @@ mod tests {
assert!(matches!(res, Err(QueryError::MissingProjectId)));
}

#[test]
fn test_generate_prefixed_id() {
let job_id = generate_prefixed_id(JOB_ID_PREFIX);
assert!(job_id.starts_with(JOB_ID_PREFIX), "{job_id:?}");
assert!(
Uuid::parse_str(&job_id[JOB_ID_PREFIX.len()..]).is_ok(),
"{job_id:?}"
);

let req_id = generate_prefixed_id(QUERY_REQUEST_ID_PREFIX);
assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX), "{req_id:?}");
assert!(req_id.len() <= 36, "{req_id:?}"); // bigquery limits request id to 36 characters
assert!(
Uuid::parse_str(&req_id[QUERY_REQUEST_ID_PREFIX.len()..]).is_ok(),
"{req_id:?}"
);
}

#[test]
fn test_generate_job_reference() {
let job_ref = generate_job_reference("my-project", "us-central1");
assert_eq!(job_ref.project_id, "my-project");
assert!(job_ref.job_id.starts_with(JOB_ID_PREFIX), "{job_ref:?}");
assert_eq!(job_ref.location.as_deref(), Some("us-central1"));
}
Comment thread
alvarowolfx marked this conversation as resolved.

#[tokio::test]
async fn test_run_jobs_insert() -> TestResult {
let mut mock = MockJobService::new();
mock.expect_insert_job().returning(|_, _| {
mock.expect_insert_job().returning(|req, _| {
let job_ref = req.job.as_ref().unwrap().job_reference.as_ref().unwrap();
assert!(job_ref.job_id.starts_with(JOB_ID_PREFIX), "{job_ref:?}");
let job_ref = JobReference::new()
.set_job_id("test-job")
.set_project_id("my-project");
Expand All @@ -167,8 +223,12 @@ mod tests {
#[tokio::test]
async fn test_run_jobs_query() -> TestResult {
let mut mock = MockJobService::new();
mock.expect_query()
.returning(move |_, _| Ok(Response::from(QueryResponse::new())));
mock.expect_query().returning(move |req, _| {
let req_id = &req.query_request.as_ref().unwrap().request_id;
assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX), "{req_id:?}");
assert!(req_id.len() <= 36, "{req_id:?}"); // bigquery limits request id to 36 characters
Ok(Response::from(QueryResponse::new()))
});
let job_service = create_job_service(mock);
let run_query =
RunQuery::new(job_service, "SELECT 1".to_string()).with_project_id("my-project");
Expand Down
Loading