Class StreamWriter (3.4.0)

public class StreamWriter implements AutoCloseable

A BigQuery Stream Writer that can be used to write data into BigQuery Table.

TODO: Support batching.

Inheritance

java.lang.Object > StreamWriter

Implements

AutoCloseable

Static Methods

getApiMaxRequestBytes()

public static long getApiMaxRequestBytes()

The maximum size of one request. Defined by the API.

Returns
TypeDescription
long

getDefaultStreamName(TableName tableName)

public static String getDefaultStreamName(TableName tableName)
Parameter
NameDescription
tableNameTableName
Returns
TypeDescription
String

the default stream name associated with tableName

newBuilder(String streamName)

public static StreamWriter.Builder newBuilder(String streamName)

Constructs a new StreamWriter.Builder using the given stream.

Parameter
NameDescription
streamNameString
Returns
TypeDescription
StreamWriter.Builder

newBuilder(String streamName, BigQueryWriteClient client)

public static StreamWriter.Builder newBuilder(String streamName, BigQueryWriteClient client)

Constructs a new StreamWriter.Builder using the given stream and client.

Parameters
NameDescription
streamNameString
clientBigQueryWriteClient
Returns
TypeDescription
StreamWriter.Builder

setMaxRequestCallbackWaitTime(Duration waitTime)

public static void setMaxRequestCallbackWaitTime(Duration waitTime)

Sets the maximum time a request is allowed to be waiting in request waiting queue. Under very low chance, it's possible for append request to be waiting indefintely for request callback when Google networking SDK does not detect the networking breakage. The default timeout is 15 minutes. We are investigating the root cause for callback not triggered by networking SDK.

Parameter
NameDescription
waitTimeDuration

Methods

append(ProtoRows rows)

public ApiFuture<AppendRowsResponse> append(ProtoRows rows)

Schedules the writing of rows at the end of current stream.

Parameter
NameDescription
rowsProtoRows

the rows in serialized format to write to BigQuery.

Returns
TypeDescription
ApiFuture<AppendRowsResponse>

the append response wrapped in a future.

append(ProtoRows rows, long offset)

public ApiFuture<AppendRowsResponse> append(ProtoRows rows, long offset)

Schedules the writing of rows at given offset.

Example of writing rows with specific offset.


 ApiFuture<AppendRowsResponse> future = writer.append(rows, 0);
 ApiFutures.addCallback(future, new ApiFutureCallback<AppendRowsResponse>() {
   public void onSuccess(AppendRowsResponse response) {
     if (!response.hasError()) {
       System.out.println("written with offset: " + response.getAppendResult().getOffset());
     } else {
       System.out.println("received an in stream error: " + response.getError().toString());
     }
   }

   public void onFailure(Throwable t) {
     System.out.println("failed to write: " + t);
   }
 }, MoreExecutors.directExecutor());
 
Parameters
NameDescription
rowsProtoRows

the rows in serialized format to write to BigQuery.

offsetlong

the offset of the first row. Provide -1 to write at the current end of stream.

Returns
TypeDescription
ApiFuture<AppendRowsResponse>

the append response wrapped in a future.

close()

public void close()

Close the stream writer. Shut down all resources.

getInflightWaitSeconds()

public long getInflightWaitSeconds()

Returns the wait of a request in Client side before sending to the Server. Request could wait in Client because it reached the client side inflight request limit (adjustable when constructing the StreamWriter). The value is the wait time for the last sent request. A constant high wait value indicates a need for more throughput, you can create a new Stream for to increase the throughput in exclusive stream case, or create a new Writer in the default stream case.

Returns
TypeDescription
long

getLocation()

public String getLocation()
Returns
TypeDescription
String

the location of the destination.

getMissingValueInterpretationMap()

public Map<String,AppendRowsRequest.MissingValueInterpretation> getMissingValueInterpretationMap()
Returns
TypeDescription
Map<String,MissingValueInterpretation>

the missing value interpretation map used for the writer.

getProtoSchema()

public ProtoSchema getProtoSchema()
Returns
TypeDescription
ProtoSchema

the passed in user schema.

getStreamName()

public String getStreamName()
Returns
TypeDescription
String

name of the Stream that this writer is working on.

getUpdatedSchema()

public synchronized TableSchema getUpdatedSchema()

Thread-safe getter of updated TableSchema.

This will return the updated schema only when the creation timestamp of this writer is older than the updated schema.

Returns
TypeDescription
TableSchema

getWriterId()

public String getWriterId()
Returns
TypeDescription
String

a unique Id for the writer.

isClosed()

public boolean isClosed()
Returns
TypeDescription
boolean

if a stream writer can no longer be used for writing. It is due to either the StreamWriter is explicitly closed or the underlying connection is broken when connection pool is not used. Client should recreate StreamWriter in this case.

isUserClosed()

public boolean isUserClosed()
Returns
TypeDescription
boolean

if user explicitly closed the writer.

setMissingValueInterpretationMap(Map<String,AppendRowsRequest.MissingValueInterpretation> missingValueInterpretationMap)

public void setMissingValueInterpretationMap(Map<String,AppendRowsRequest.MissingValueInterpretation> missingValueInterpretationMap)

Sets the missing value interpretation map for the stream writer. The input missingValueInterpretationMap is used for all write requests unless otherwise changed.

Parameter
NameDescription
missingValueInterpretationMapMap<String,MissingValueInterpretation>

the missing value interpretation map used by stream writer.