Class BigQuerySink
java.lang.Object
io.github.flink.gcp.connector.bigquery.sink.BigQuerySink
Entry point for building a BigQuery sink.
The sink exposes a single builder-based API and dispatches, at job-graph construction time, to
one of several write-method implementations, in the spirit of Apache Beam's BigQueryIO:
WriteMethod.STORAGE_API_AT_LEAST_ONCE— BigQuery Storage Write API default stream, at-least-once semantics, supporting dynamic per-record table destinationsWriteMethod.STORAGE_API_EXACTLY_ONCE— BigQuery Storage Write API buffered streams with two-phase commit, exactly-once semanticsWriteMethod.FILE_LOADS— files staged on Cloud Storage followed by BigQuery load jobs, exactly-once in batch and checkpointed streaming execution
Those semantics assume the default FailureHandler.failJob() policy. Under a dropping
policy configured through BigQuerySinkBuilder.failureHandler(FailureHandler), they cover
every record except those handed to that handler, which are never written at all. A record the
serializer skips by returning null is written nowhere either, under any policy.
Write methods that are not implemented yet are rejected by BigQuerySinkBuilder.build()
with an UnsupportedOperationException.
Example:
Sink<MyEvent> sink =
BigQuerySink.<MyEvent>builder()
.writeMethod(WriteMethod.STORAGE_API_AT_LEAST_ONCE)
.destinationResolver(
(e, ctx) ->
TableDestination.of(
"my-project", "my_dataset", e.tableName()))
.serializer(new MyEventProtoSerializer())
.build();
-
Method Summary
Modifier and TypeMethodDescriptionstatic <T> BigQuerySinkBuilder<T>builder()Creates a newBigQuerySinkBuilder.
-
Method Details
-
builder
Creates a newBigQuerySinkBuilder.- Type Parameters:
T- type of the records written by the sink- Returns:
- a new builder
-