pub struct WriterBuilder<S> { /* private fields */ }Expand description
A builder to create a stream writer.
Implementations§
Source§impl WriterBuilder<DefaultStream>
impl WriterBuilder<DefaultStream>
Sourcepub fn with_multiplexing(self, enable: bool) -> Self
pub fn with_multiplexing(self, enable: bool) -> Self
Enable multiplexing.
Set this option to use the client’s shared stream pool.
This option only applies to the default stream.
Note that until background stream rebalancing (#6866) is implemented, stream assignment and pool scale-up only occur when a writer is built. If multiple writers are created all at once before writes are in flight, they will all share the same initial stream. To scale across multiple streams, stagger writer creation so writes are in flight when new writers are built.
§Example
let writer = client
.open_default_stream("projects/my-project/datasets/my_dataset/tables/my_table")
.with_multiplexing(true)
.build_arrow(schema())
.await?;Sourcepub fn with_retry_policy<V: Into<RetryPolicyArg>>(self, v: V) -> Self
pub fn with_retry_policy<V: Into<RetryPolicyArg>>(self, v: V) -> Self
Configure the retry policy.
The client libraries can automatically retry operations that fail. The retry policy controls what errors are considered retryable, sets limits on the number of attempts or the time trying to make attempts.
§Example
use google_cloud_bigquery::write::retry_policy::RetryableErrors;
use google_cloud_gax::retry_policy::RetryPolicyExt;
let writer = client
.open_default_stream("projects/my-project/datasets/my_dataset/tables/my_table")
.with_retry_policy(RetryableErrors.with_attempt_limit(3))
.build_arrow(schema())
.await?;Sourcepub fn with_backoff_policy<V: Into<BackoffPolicyArg>>(self, v: V) -> Self
pub fn with_backoff_policy<V: Into<BackoffPolicyArg>>(self, v: V) -> Self
Configure the retry backoff policy.
The client libraries can automatically retry operations that fail. The backoff policy controls how long to wait in between retry attempts.
§Example
use google_cloud_gax::exponential_backoff::ExponentialBackoff;
let policy = ExponentialBackoff::default();
let writer = client
.open_default_stream("projects/my-project/datasets/my_dataset/tables/my_table")
.with_backoff_policy(policy)
.build_arrow(schema())
.await?;Sourcepub fn with_attempt_timeout(self, v: Duration) -> Self
pub fn with_attempt_timeout(self, v: Duration) -> Self
Configure the timeout for a single write attempt.
Without this limit, a write can block forever if the service accepts the stream but never responds. On a timeout, the client abandons the stream and the retry policy decides whether to make another attempt.
§Example
use std::time::Duration;
let writer = client
.open_default_stream("projects/my-project/datasets/my_dataset/tables/my_table")
.with_attempt_timeout(Duration::from_secs(10))
.build_arrow(schema())
.await?;Source§impl<S: Stream> WriterBuilder<S>
impl<S: Stream> WriterBuilder<S>
Sourcepub async fn build_arrow<W>(
self,
schema: ArrowSchema,
) -> Result<W, WriterBuilderError>
pub async fn build_arrow<W>( self, schema: ArrowSchema, ) -> Result<W, WriterBuilderError>
Consumes the builder and creates a writer using Arrow as the data format.
Returns the writer W corresponding to the stream type S:
DefaultStream->DefaultWriter<Arrow>PendingStream->PendingWriter<Arrow>CommittedStream->CommittedWriter<Arrow>BufferedStream->BufferedWriter<Arrow>
§Example
let writer = client
.open_default_stream("projects/my-project/datasets/my_dataset/tables/my_table")
.build_arrow(ArrowSchema::new())
.await?;Trait Implementations§
Source§impl<S: Clone> Clone for WriterBuilder<S>
impl<S: Clone> Clone for WriterBuilder<S>
Auto Trait Implementations§
impl<S> !RefUnwindSafe for WriterBuilder<S>
impl<S> !UnwindSafe for WriterBuilder<S>
impl<S> Freeze for WriterBuilder<S>where
PhantomData<S>: Freeze,
impl<S> Send for WriterBuilder<S>where
PhantomData<S>: Send,
impl<S> Sync for WriterBuilder<S>where
PhantomData<S>: Sync,
impl<S> Unpin for WriterBuilder<S>where
PhantomData<S>: Unpin,
impl<S> UnsafeUnpin for WriterBuilder<S>where
PhantomData<S>: UnsafeUnpin,
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
Source§impl<T> FutureExt for T
impl<T> FutureExt for T
Source§fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
Source§fn with_current_context(self) -> WithContext<Self> ⓘ
fn with_current_context(self) -> WithContext<Self> ⓘ
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoRequest<T> for T
impl<T> IntoRequest<T> for T
Source§fn into_request(self) -> Request<T>
fn into_request(self) -> Request<T>
T in a tonic::Request