Problem
AWS Kinesis Firehose delivers Database Activity Streams to S3. Depending on Firehose buffer size and record grouping configuration, Firehose frequently buffers multiple records into a single S3 object:
- Newline-delimited JSON (NDJSON): Multiple JSON objects separated by
\n
- Concatenated JSON: Multiple adjacent JSON objects without delimiters (e.g.
{"type":...}{"type":...})
In main.go, json.Unmarshal is used directly on the whole object body:
var payload DASPayload
if err := json.Unmarshal(bodyBytes, &payload); err != nil {
return fmt.Errorf("failed to parse DAS payload: %w", err)
}
If Firehose buffers more than one record into the S3 object, json.Unmarshal fails with:
invalid character '{' after top-level value
This causes the entire batch to fail and unretried records to be stuck in the queue.
Proposed Solution
Use a streaming json.Decoder to iterate over all top-level DASPayload records in the S3 object:
decoder := json.NewDecoder(bytes.NewReader(bodyBytes))
for decoder.More() {
var payload DASPayload
if err := decoder.Decode(&payload); err != nil {
return fmt.Errorf("failed to decode DAS payload: %w", err)
}
// Process payload...
}
This handles both single JSON objects, NDJSON, and concatenated JSON objects seamlessly.
Acceptance Criteria
Problem
AWS Kinesis Firehose delivers Database Activity Streams to S3. Depending on Firehose buffer size and record grouping configuration, Firehose frequently buffers multiple records into a single S3 object:
\n{"type":...}{"type":...})In main.go,
json.Unmarshalis used directly on the whole object body:If Firehose buffers more than one record into the S3 object,
json.Unmarshalfails with:This causes the entire batch to fail and unretried records to be stuck in the queue.
Proposed Solution
Use a streaming
json.Decoderto iterate over all top-level DASPayload records in the S3 object:This handles both single JSON objects, NDJSON, and concatenated JSON objects seamlessly.
Acceptance Criteria