Skip to content
Merged
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
111 changes: 56 additions & 55 deletions out_writeapi.go
Original file line number Diff line number Diff line change
Expand Up @@ -58,18 +58,16 @@ type outputConfig struct {
managedStreamSlice *[]*streamConfig
client ManagedWriterClient
maxChunkSize int
mutex sync.Mutex
streamMu sync.Mutex // guards managedStreamSlice and its elements (streamConfig)
exactlyOnce bool
requestCountThreshold int
numRetries int
timestampFields []string
fieldCache fieldLookupCache
flushTimeout time.Duration
binaryData [][]byte // reused across flushes to avoid per-flush allocation
}

var (
ms_ctx = context.Background()
configMap = make(map[int]*outputConfig)
configID atomic.Int64
)
Expand All @@ -85,7 +83,7 @@ const (
minQueueRequests = 10 // minimum Max_Queue_Requests accepted (prevents starvation during scaling)
dateTimeDefault = true
maxUnixSeconds = 4102444800 // 2100-01-01 00:00:00 UTC; values above this are treated as already in microseconds
flushTimeoutSecDefault = 30 // default Flush_Timeout_Sec: 30 seconds per flush call
flushTimeoutSecDefault = 10 // default Flush_Timeout_Sec: 10 seconds per flush call
)

// This function mangles the top-level and complex (struct) BigQuery schema to convert NUMERIC, BIGNUMERIC, DATETIME, TIME, and JSON fields to STRING.
Expand Down Expand Up @@ -381,10 +379,12 @@ func rebuildPredicate(err error) bool {
return false
}

// This function sends and checks the responses for data through a committed stream with exactly once functionality
// This function sends and checks the responses for data through a committed stream with exactly once functionality.
// On success the stream's offsetCounter is incremented by len(data) within the same streamMu critical section
// as the append, preventing TOCTOU races when workers > 1.
func sendRequestExactlyOnce(ctx context.Context, data [][]byte, config *outputConfig, streamIndex int) error {
config.mutex.Lock()
defer config.mutex.Unlock()
config.streamMu.Lock()
defer config.streamMu.Unlock()

currStream := (*config.managedStreamSlice)[streamIndex]

Expand All @@ -397,27 +397,37 @@ func sendRequestExactlyOnce(ctx context.Context, data [][]byte, config *outputCo
if err != nil {
return err
}
// Increment the offset inside the same lock so no other worker can read a stale offset value.
currStream.offsetCounter += int64(len(data))
return nil
}

// This function enables synchronous retries and rebuilding a valid stream based on the server response
// This function enables synchronous retries and rebuilding a valid stream based on the server response.
// The streamMu is held while closing and rebuilding the stream (rebuildPredicate path) so that no other
// worker (workers > 1 scenario) can use the stream concurrently during the teardown/rebuild.
// Note: this means network I/O (Finalize, Close, buildStream) runs under the lock on that path,
// which is a known trade-off; see issue #12 in docs/code-analysis.md.
func sendRequestRetries(ctx context.Context, data [][]byte, config *outputConfig, streamIndex int) error {
retryer := newStatelessRetryer(config.numRetries)
attempt := 0
currStream := (*config.managedStreamSlice)[streamIndex]
for {
err := sendRequestExactlyOnce(ctx, data, config, streamIndex)
if err == nil {
break
}
// Unsuccessful data append
if rebuildPredicate(err) {
// Hold streamMu while closing and rebuilding the stream so no concurrent
// worker (workers > 1) can observe a partially-closed or replaced stream.
config.streamMu.Lock()
currStream := (*config.managedStreamSlice)[streamIndex]
currStream.managedstream.Finalize(ctx)
currStream.managedstream.Close()
// Rebuild stream
err := buildStream(ctx, config, streamIndex)
if err != nil {
return err
buildErr := buildStream(ctx, config, streamIndex)
config.streamMu.Unlock()
if buildErr != nil {
return buildErr
}
// Retry sending data without incrementing number of attempts or waiting between attempts
} else {
Expand All @@ -436,8 +446,8 @@ func sendRequestRetries(ctx context.Context, data [][]byte, config *outputConfig

// This function sends data and appends the responses to a queue to be checked asynchronously through a default stream with at least once functionality
func sendRequestDefault(ctx context.Context, data [][]byte, config *outputConfig, streamIndex int) error {
config.mutex.Lock()
defer config.mutex.Unlock()
config.streamMu.Lock()
defer config.streamMu.Unlock()
currStream := (*config.managedStreamSlice)[streamIndex]

appendResult, err := currStream.managedstream.AppendRows(ctx, data)
Expand Down Expand Up @@ -498,8 +508,8 @@ var setThreshold = func(maxQueueSize int) int {
// This function check whether there is room for scaling and the scales the number of stream dynamically depending on if it
// Detects back pressure from the queue
func createNewStreamDynamicScaling(ctx context.Context, config *outputConfig) {
config.mutex.Lock()
defer config.mutex.Unlock()
config.streamMu.Lock()
defer config.streamMu.Unlock()
if len(*config.managedStreamSlice) < maxNumStreamsPerInstance {
// Gets stream with least values in queue
mostEfficient := getLeastLoadedStream(config.managedStreamSlice)
Expand Down Expand Up @@ -598,15 +608,15 @@ var getFLBPluginContext = func(ctx unsafe.Pointer) int {
}

// Finalizes and Closes all streams in slice for a given instance
func finalizeCloseAllStreams(config *outputConfig, id int) bool {
config.mutex.Lock()
defer config.mutex.Unlock()
func finalizeCloseAllStreams(ctx context.Context, config *outputConfig, id int) bool {
config.streamMu.Lock()
defer config.streamMu.Unlock()
errFlag := false
streamSlice := config.managedStreamSlice
for i := 0; i < len(*config.managedStreamSlice); i++ {
if (*streamSlice)[i].managedstream != nil {
if config.exactlyOnce {
if _, err := (*streamSlice)[i].managedstream.Finalize(ms_ctx); err != nil {
if _, err := (*streamSlice)[i].managedstream.Finalize(ctx); err != nil {
log.Printf("Finalizing managed stream for output instance with id %d and stream index %d failed in FLBPluginExit: %s", id, i, err)
errFlag = true
}
Expand Down Expand Up @@ -686,7 +696,8 @@ func FLBPluginInit(plugin unsafe.Pointer) int {
}

// Create new client
client, err := getClient(ms_ctx, projectID)
initCtx := context.Background()
client, err := getClient(initCtx, projectID)
if err != nil {
log.Printf("Creating a new managed BigQuery Storage write client scoped to: %s failed in FLBPluginInit: %s", projectID, err)
return output.FLB_ERROR
Expand All @@ -696,7 +707,7 @@ func FLBPluginInit(plugin unsafe.Pointer) int {
tableReference := fmt.Sprintf("projects/%s/datasets/%s/tables/%s", projectID, datasetID, tableID)

// Call getDescriptors to get the message descriptor, and descriptor proto
md, descriptor, timestampFields, err := getDescriptors(ms_ctx, client, projectID, datasetID, tableID, dateTimeStringType)
md, descriptor, timestampFields, err := getDescriptors(initCtx, client, projectID, datasetID, tableID, dateTimeStringType)
if err != nil {
log.Printf("Getting message descriptor and descriptor proto for table: %s failed in FLBPluginInit: %s", tableReference, err)
return output.FLB_ERROR
Expand Down Expand Up @@ -745,7 +756,7 @@ func FLBPluginInit(plugin unsafe.Pointer) int {
}

// Create stream using NewManagedStream
err = buildStream(ms_ctx, &config, 0)
err = buildStream(initCtx, &config, 0)
if err != nil {
log.Printf("Creating a new managed stream with destination table: %s failed in FLBPluginInit: %s", tableReference, err)
return output.FLB_ERROR
Expand All @@ -766,32 +777,27 @@ func FLBPluginFlush(data unsafe.Pointer, length C.int, tag *C.char) int {
return output.FLB_OK
}

// flushChunk picks the least-loaded stream, sends binaryData, and (for
// exactly-once mode) increments the stream's offset counter by rowCount.
func flushChunk(ctx context.Context, config *outputConfig, id int, binaryData [][]byte, rowCount int64) {
config.mutex.Lock()
// flushChunk picks the least-loaded stream and sends binaryData.
// For exactly-once mode the stream's offsetCounter is incremented by len(binaryData)
// inside sendRequestExactlyOnce within the same streamMu critical section as the
// AppendRows call, preventing TOCTOU races when workers > 1.
func flushChunk(ctx context.Context, config *outputConfig, id int, binaryData [][]byte) {
config.streamMu.Lock()
streamIndex := getLeastLoadedStream(config.managedStreamSlice)
config.mutex.Unlock()
config.streamMu.Unlock()

if err := sendRequest(ctx, binaryData, config, streamIndex); err != nil {
log.Printf("Appending data for output instance with id: %d failed in FLBPluginFlushCtx: %s", id, err)
return
}
if config.exactlyOnce {
config.mutex.Lock()
(*config.managedStreamSlice)[streamIndex].offsetCounter += rowCount
config.mutex.Unlock()
}
}

// decodeAndSerializeRecords decodes Fluent Bit records from dec, serializes each
// to proto binary, and accumulates them in binaryData. When the accumulated size
// reaches config.maxChunkSize, it calls flushChunk and resets binaryData
// (preserving capacity for reuse). It returns the remaining (unsent) binaryData
// and the row count for the final chunk.
func decodeAndSerializeRecords(ctx context.Context, config *outputConfig, id int, dec *output.FLBDecoder, binaryData [][]byte) ([][]byte, int64) {
// (preserving capacity for reuse). It returns the remaining (unsent) binaryData.
func decodeAndSerializeRecords(ctx context.Context, config *outputConfig, id int, dec *output.FLBDecoder, binaryData [][]byte) [][]byte {
var currsize int
var rowCounter int64

for {
ret, _, record := output.GetRecord(dec)
Expand All @@ -808,8 +814,7 @@ func decodeAndSerializeRecords(ctx context.Context, config *outputConfig, id int
}

if (currsize + len(buf)) >= config.maxChunkSize {
flushChunk(ctx, config, id, binaryData, rowCounter)
rowCounter = 0
flushChunk(ctx, config, id, binaryData)
// Nil out sent elements so GC can reclaim them, then reset slice.
for i := range binaryData {
binaryData[i] = nil
Expand All @@ -820,10 +825,9 @@ func decodeAndSerializeRecords(ctx context.Context, config *outputConfig, id int
binaryData = append(binaryData, buf)
// Include the protobuf overhead in the size estimate.
currsize += (len(buf) + 2)
rowCounter++
}

return binaryData, rowCounter
return binaryData
}

//export FLBPluginFlushCtx
Expand All @@ -839,24 +843,18 @@ func FLBPluginFlushCtx(ctx, data unsafe.Pointer, length C.int, tag *C.char) int
defer flushCancel()

// Drain ready responses and check whether a new stream should be created.
checkAllStreamResponses(flushCtx, &config.managedStreamSlice, false, &config.mutex, config.exactlyOnce, id)
checkAllStreamResponses(flushCtx, &config.managedStreamSlice, false, &config.streamMu, config.exactlyOnce, id)
createNewStreamDynamicScaling(flushCtx, config)

// Reuse binaryData slice across flushes to avoid per-flush allocation.
// Nil out retained elements before reset so GC can reclaim previous batch's bytes.
for i := range config.binaryData {
config.binaryData[i] = nil
}
binaryData := config.binaryData[:0]
// binaryData is a flush-local buffer. Each worker (workers > 1) gets its own
// stack frame, so there is no shared state and no mutex needed here.
binaryData := make([][]byte, 0, 64)

// Decode and serialize all records; mid-size chunks are sent inside the helper.
binaryData, rowCounter := decodeAndSerializeRecords(flushCtx, config, id, newDecoder(data, int(length)), binaryData)
binaryData = decodeAndSerializeRecords(flushCtx, config, id, newDecoder(data, int(length)), binaryData)

// Send the final (possibly partial) chunk.
flushChunk(flushCtx, config, id, binaryData, rowCounter)

// Write back the grown slice so its capacity is reused on the next flush.
config.binaryData = binaryData
flushChunk(flushCtx, config, id, binaryData)

return output.FLB_OK
}
Expand All @@ -881,8 +879,11 @@ func FLBPluginExitCtx(ctx unsafe.Pointer) int {

// Drain all pending responses before closing streams to avoid data loss.
// waitForResponse=true blocks until every in-flight AppendRows result is received.
checkAllStreamResponses(ms_ctx, &config.managedStreamSlice, true, &config.mutex, config.exactlyOnce, id)
errFlag := finalizeCloseAllStreams(config, id)
// Use a timeout context (same duration as flush) so exit cannot hang indefinitely.
exitCtx, exitCancel := context.WithTimeout(context.Background(), config.flushTimeout)
defer exitCancel()
checkAllStreamResponses(exitCtx, &config.managedStreamSlice, true, &config.streamMu, config.exactlyOnce, id)
errFlag := finalizeCloseAllStreams(exitCtx, config, id)

if config.client != nil {
if err := config.client.Close(); err != nil {
Expand Down