Skip to content

Commit e737867

Browse files
committed
parquet_encode: add default_timestamp_unit field
Adds a new `default_timestamp_unit` configuration field to the `parquet_encode` processor, accepting `NANOSECOND` (default, preserves existing behaviour), `MICROSECOND`, or `MILLISECOND`. The unit is applied to both the static schema path and the dynamic `schema_metadata` path (used by CDC inputs such as `mysql_cdc`). `TIMESTAMP(NANOS)` is not readable by Apache Spark / Databricks, AWS Athena or DuckDB; this field unblocks those consumers without requiring a pre-encoding transform. MySQL sources additionally cannot exceed microsecond precision, so `MICROSECOND` is lossless for CDC pipelines. Resolves the long-standing TODO referenced at #3570.
1 parent cb78ded commit e737867

4 files changed

Lines changed: 160 additions & 21 deletions

File tree

CHANGELOG.md

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,12 @@ Changelog
33

44
All notable changes to this project will be documented in this file.
55

6+
## Unreleased
7+
8+
### Added
9+
10+
- parquet_encode: Added `default_timestamp_unit` field (values `NANOSECOND`, `MICROSECOND`, `MILLISECOND`) controlling the precision of TIMESTAMP logical types. Default remains `NANOSECOND` for backwards compatibility. Use `MICROSECOND` when writing files for Apache Spark/Databricks, AWS Athena or DuckDB, which do not support `TIMESTAMP(NANOS)`. ([#3570](https://github.com/redpanda-data/connect/issues/3570))
11+
612
## 4.88.0 - 2026-04-16
713

814
### Added

docs/modules/components/pages/processors/parquet_encode.adoc

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,7 @@ parquet_encode:
5656
schema_metadata: ""
5757
default_compression: uncompressed
5858
default_encoding: DELTA_LENGTH_BYTE_ARRAY
59+
default_timestamp_unit: NANOSECOND
5960
```
6061
6162
--
@@ -216,4 +217,19 @@ Options:
216217
, `PLAIN`
217218
.
218219
220+
=== `default_timestamp_unit`
221+
222+
The precision used when encoding TIMESTAMP logical types. The default `NANOSECOND` matches historical behaviour, but `TIMESTAMP(NANOS)` is not readable by Apache Spark (Databricks), AWS Athena or DuckDB; set this to `MICROSECOND` (or `MILLISECOND`) when writing Parquet files intended for consumption by those engines.
223+
224+
225+
*Type*: `string`
226+
227+
*Default*: `"NANOSECOND"`
228+
229+
Options:
230+
`NANOSECOND`
231+
, `MICROSECOND`
232+
, `MILLISECOND`
233+
.
234+
219235

internal/impl/parquet/processor_encode.go

Lines changed: 44 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,12 @@ func parquetEncodeProcessorConfig() *service.ConfigSpec {
4949
Default("DELTA_LENGTH_BYTE_ARRAY").
5050
Advanced().
5151
Version("4.11.0"),
52+
service.NewStringEnumField("default_timestamp_unit",
53+
"NANOSECOND", "MICROSECOND", "MILLISECOND",
54+
).
55+
Description("The precision used when encoding TIMESTAMP logical types. The default `NANOSECOND` matches historical behaviour, but `TIMESTAMP(NANOS)` is not readable by Apache Spark (Databricks), AWS Athena or DuckDB; set this to `MICROSECOND` (or `MILLISECOND`) when writing Parquet files intended for consumption by those engines.").
56+
Default("NANOSECOND").
57+
Advanced(),
5258
).
5359
Description(`
5460
This processor uses https://github.com/parquet-go/parquet-go[https://github.com/parquet-go/parquet-go^], which is itself experimental. Therefore changes could be made into how this processor functions outside of major version releases.
@@ -119,7 +125,19 @@ var plainEncodingFn encodingFn = func(n parquet.Node) parquet.Node {
119125
return parquet.Encoded(n, &parquet.Plain)
120126
}
121127

122-
func parquetGroupFromConfig(columnConfs []*service.ParsedConfig, encodingFn encodingFn) (parquet.Group, error) {
128+
func parseTimestampUnit(s string) (parquet.TimeUnit, error) {
129+
switch s {
130+
case "NANOSECOND":
131+
return parquet.Nanosecond, nil
132+
case "MICROSECOND":
133+
return parquet.Microsecond, nil
134+
case "MILLISECOND":
135+
return parquet.Millisecond, nil
136+
}
137+
return nil, fmt.Errorf("unknown timestamp unit %q (expected NANOSECOND, MICROSECOND or MILLISECOND)", s)
138+
}
139+
140+
func parquetGroupFromConfig(columnConfs []*service.ParsedConfig, encodingFn encodingFn, tsUnit parquet.TimeUnit) (parquet.Group, error) {
123141
groupNode := parquet.Group{}
124142

125143
for _, colConf := range columnConfs {
@@ -131,7 +149,7 @@ func parquetGroupFromConfig(columnConfs []*service.ParsedConfig, encodingFn enco
131149
}
132150

133151
if childColumns, _ := colConf.FieldAnyList("fields"); len(childColumns) > 0 {
134-
if n, err = parquetGroupFromConfig(childColumns, encodingFn); err != nil {
152+
if n, err = parquetGroupFromConfig(childColumns, encodingFn, tsUnit); err != nil {
135153
return nil, err
136154
}
137155
} else {
@@ -155,8 +173,7 @@ func parquetGroupFromConfig(columnConfs []*service.ParsedConfig, encodingFn enco
155173
case "UTF8":
156174
n = parquet.String()
157175
case "TIMESTAMP":
158-
// TODO: add field to specify timestamp unit (https://github.com/redpanda-data/connect/issues/3570)
159-
n = parquet.Timestamp(parquet.Nanosecond)
176+
n = parquet.Timestamp(tsUnit)
160177
case "BSON":
161178
n = parquet.BSON()
162179
case "ENUM":
@@ -193,6 +210,15 @@ func parquetGroupFromConfig(columnConfs []*service.ParsedConfig, encodingFn enco
193210
//------------------------------------------------------------------------------
194211

195212
func newParquetEncodeProcessorFromConfig(conf *service.ParsedConfig, logger *service.Logger) (*parquetEncodeProcessor, error) {
213+
tsUnitStr, err := conf.FieldString("default_timestamp_unit")
214+
if err != nil {
215+
return nil, err
216+
}
217+
tsUnit, err := parseTimestampUnit(tsUnitStr)
218+
if err != nil {
219+
return nil, err
220+
}
221+
196222
var schema *parquet.Schema
197223
if conf.Contains("schema") {
198224
schemaConfs, err := conf.FieldObjectList("schema")
@@ -212,7 +238,7 @@ func newParquetEncodeProcessorFromConfig(conf *service.ParsedConfig, logger *ser
212238
encoding = defaultEncodingFn
213239
}
214240

215-
node, err := parquetGroupFromConfig(schemaConfs, encoding)
241+
node, err := parquetGroupFromConfig(schemaConfs, encoding, tsUnit)
216242
if err != nil {
217243
return nil, err
218244
}
@@ -250,22 +276,24 @@ func newParquetEncodeProcessorFromConfig(conf *service.ParsedConfig, logger *ser
250276
default:
251277
return nil, fmt.Errorf("default_compression type %v not recognised", compressStr)
252278
}
253-
return newParquetEncodeProcessor(logger, schema, schemaMeta, compressDefault)
279+
return newParquetEncodeProcessor(logger, schema, schemaMeta, compressDefault, tsUnit)
254280
}
255281

256282
type parquetEncodeProcessor struct {
257283
logger *service.Logger
258284
schema *parquet.Schema
259285
schemaMeta string
260286
compressionType compress.Codec
287+
timestampUnit parquet.TimeUnit
261288
}
262289

263-
func newParquetEncodeProcessor(logger *service.Logger, schema *parquet.Schema, schemaMeta string, compressionType compress.Codec) (*parquetEncodeProcessor, error) {
290+
func newParquetEncodeProcessor(logger *service.Logger, schema *parquet.Schema, schemaMeta string, compressionType compress.Codec, timestampUnit parquet.TimeUnit) (*parquetEncodeProcessor, error) {
264291
s := &parquetEncodeProcessor{
265292
logger: logger,
266293
schema: schema,
267294
schemaMeta: schemaMeta,
268295
compressionType: compressionType,
296+
timestampUnit: timestampUnit,
269297
}
270298
return s, nil
271299
}
@@ -305,7 +333,7 @@ func (s *parquetEncodeProcessor) ProcessBatch(_ context.Context, batch service.M
305333
}
306334

307335
var err error
308-
if schema, err = parquetSchemaFromCommon(metaAny); err != nil {
336+
if schema, err = parquetSchemaFromCommon(metaAny, s.timestampUnit); err != nil {
309337
return nil, err
310338
}
311339
}
@@ -348,7 +376,7 @@ func (*parquetEncodeProcessor) Close(context.Context) error {
348376
return nil
349377
}
350378

351-
func parquetNodeFromCommonField(field schema.Common) (parquet.Node, error) {
379+
func parquetNodeFromCommonField(field schema.Common, tsUnit parquet.TimeUnit) (parquet.Node, error) {
352380
var n parquet.Node
353381

354382
switch field.Type {
@@ -365,8 +393,7 @@ func parquetNodeFromCommonField(field schema.Common) (parquet.Node, error) {
365393
case schema.String:
366394
n = parquet.String()
367395
case schema.Timestamp:
368-
// TODO: add field to specify timestamp unit (https://github.com/redpanda-data/connect/issues/3570)
369-
n = parquet.Timestamp(parquet.Nanosecond)
396+
n = parquet.Timestamp(tsUnit)
370397
case schema.ByteArray:
371398
n = parquet.Leaf(parquet.ByteArrayType)
372399
case schema.Array:
@@ -375,7 +402,7 @@ func parquetNodeFromCommonField(field schema.Common) (parquet.Node, error) {
375402
}
376403

377404
var err error
378-
if n, err = parquetNodeFromCommonField(field.Children[0]); err != nil {
405+
if n, err = parquetNodeFromCommonField(field.Children[0], tsUnit); err != nil {
379406
return nil, err
380407
}
381408
n = parquet.Repeated(n)
@@ -386,7 +413,7 @@ func parquetNodeFromCommonField(field schema.Common) (parquet.Node, error) {
386413
}
387414

388415
var err error
389-
if n, err = parquetGroupFromCommonFields(field.Children); err != nil {
416+
if n, err = parquetGroupFromCommonFields(field.Children, tsUnit); err != nil {
390417
return nil, err
391418
}
392419

@@ -403,11 +430,11 @@ func parquetNodeFromCommonField(field schema.Common) (parquet.Node, error) {
403430
return n, nil
404431
}
405432

406-
func parquetGroupFromCommonFields(fields []schema.Common) (parquet.Group, error) {
433+
func parquetGroupFromCommonFields(fields []schema.Common, tsUnit parquet.TimeUnit) (parquet.Group, error) {
407434
g := parquet.Group{}
408435

409436
for _, f := range fields {
410-
n, err := parquetNodeFromCommonField(f)
437+
n, err := parquetNodeFromCommonField(f, tsUnit)
411438
if err != nil {
412439
return nil, err
413440
}
@@ -417,7 +444,7 @@ func parquetGroupFromCommonFields(fields []schema.Common) (parquet.Group, error)
417444
return g, nil
418445
}
419446

420-
func parquetSchemaFromCommon(a any) (*parquet.Schema, error) {
447+
func parquetSchemaFromCommon(a any, tsUnit parquet.TimeUnit) (*parquet.Schema, error) {
421448
commonSchema, err := schema.ParseFromAny(a)
422449
if err != nil {
423450
return nil, err
@@ -431,7 +458,7 @@ func parquetSchemaFromCommon(a any) (*parquet.Schema, error) {
431458
return nil, fmt.Errorf("source schema must have at least one field, got %v", len(commonSchema.Children))
432459
}
433460

434-
groupNode, err := parquetGroupFromCommonFields(commonSchema.Children)
461+
groupNode, err := parquetGroupFromCommonFields(commonSchema.Children, tsUnit)
435462
if err != nil {
436463
return nil, err
437464
}

internal/impl/parquet/processor_encode_test.go

Lines changed: 94 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -332,7 +332,7 @@ func TestParquetEncodeProcessor(t *testing.T) {
332332
expectedDataBytes, err := json.Marshal(test.input)
333333
require.NoError(t, err)
334334

335-
reader, err := newParquetEncodeProcessor(nil, testPMSchema(), "", &parquet.Uncompressed)
335+
reader, err := newParquetEncodeProcessor(nil, testPMSchema(), "", &parquet.Uncompressed, parquet.Nanosecond)
336336
require.NoError(t, err)
337337

338338
readerResBatches, err := reader.ProcessBatch(t.Context(), service.MessageBatch{
@@ -379,7 +379,7 @@ func TestParquetEncodeProcessor(t *testing.T) {
379379
inBatch = append(inBatch, service.NewMessage(dataBytes))
380380
}
381381

382-
reader, err := newParquetEncodeProcessor(nil, testPMSchema(), "", &parquet.Uncompressed)
382+
reader, err := newParquetEncodeProcessor(nil, testPMSchema(), "", &parquet.Uncompressed, parquet.Nanosecond)
383383
require.NoError(t, err)
384384

385385
readerResBatches, err := reader.ProcessBatch(t.Context(), inBatch)
@@ -578,7 +578,7 @@ func TestParquetEncodeDynamicSchemaProcessor(t *testing.T) {
578578

579579
inBatch[0].MetaSetMut("foobar", commonSchema.ToAny())
580580

581-
reader, err := newParquetEncodeProcessor(nil, nil, "foobar", &parquet.Uncompressed)
581+
reader, err := newParquetEncodeProcessor(nil, nil, "foobar", &parquet.Uncompressed, parquet.Nanosecond)
582582
require.NoError(t, err)
583583

584584
readerResBatches, err := reader.ProcessBatch(t.Context(), inBatch)
@@ -638,7 +638,7 @@ func TestParquetEncodeDynamicSchemaAnyFieldError(t *testing.T) {
638638
}
639639
inBatch[0].MetaSetMut("schema", commonSchema.ToAny())
640640

641-
proc, err := newParquetEncodeProcessor(nil, nil, "schema", &parquet.Uncompressed)
641+
proc, err := newParquetEncodeProcessor(nil, nil, "schema", &parquet.Uncompressed, parquet.Nanosecond)
642642
require.NoError(t, err)
643643

644644
_, err = proc.ProcessBatch(t.Context(), inBatch)
@@ -647,6 +647,96 @@ func TestParquetEncodeDynamicSchemaAnyFieldError(t *testing.T) {
647647
assert.Contains(t, err.Error(), "ANY")
648648
}
649649

650+
func TestParquetEncodeTimestampUnit(t *testing.T) {
651+
tests := []struct {
652+
name string
653+
unitConfig string
654+
expectedUnit parquet.TimeUnit
655+
expectedSchema string
656+
}{
657+
{name: "default is nanosecond", unitConfig: "", expectedUnit: parquet.Nanosecond, expectedSchema: "unit=NANOS"},
658+
{name: "microsecond", unitConfig: "default_timestamp_unit: MICROSECOND", expectedUnit: parquet.Microsecond, expectedSchema: "unit=MICROS"},
659+
{name: "millisecond", unitConfig: "default_timestamp_unit: MILLISECOND", expectedUnit: parquet.Millisecond, expectedSchema: "unit=MILLIS"},
660+
{name: "explicit nanosecond", unitConfig: "default_timestamp_unit: NANOSECOND", expectedUnit: parquet.Nanosecond, expectedSchema: "unit=NANOS"},
661+
}
662+
663+
for _, test := range tests {
664+
t.Run(test.name, func(t *testing.T) {
665+
configYAML := fmt.Sprintf(`
666+
schema:
667+
- { name: id, type: INT64 }
668+
- { name: ts, type: TIMESTAMP }
669+
%s
670+
`, test.unitConfig)
671+
encodeConf, err := parquetEncodeProcessorConfig().ParseYAML(configYAML, nil)
672+
require.NoError(t, err)
673+
674+
encodeProc, err := newParquetEncodeProcessorFromConfig(encodeConf, nil)
675+
require.NoError(t, err)
676+
require.Equal(t, test.expectedUnit, encodeProc.timestampUnit)
677+
678+
batches, err := encodeProc.ProcessBatch(t.Context(), service.MessageBatch{
679+
service.NewMessage([]byte(`{"id":1,"ts":"2026-04-17T12:00:00Z"}`)),
680+
})
681+
require.NoError(t, err)
682+
require.Len(t, batches, 1)
683+
require.Len(t, batches[0], 1)
684+
685+
pqBytes, err := batches[0][0].AsBytes()
686+
require.NoError(t, err)
687+
688+
pqFile, err := parquet.OpenFile(bytes.NewReader(pqBytes), int64(len(pqBytes)))
689+
require.NoError(t, err)
690+
assert.Contains(t, pqFile.Schema().String(), test.expectedSchema)
691+
})
692+
}
693+
}
694+
695+
func TestParquetEncodeTimestampUnitDynamicSchema(t *testing.T) {
696+
encodeConf, err := parquetEncodeProcessorConfig().ParseYAML(`
697+
schema_metadata: benthos_schema
698+
default_timestamp_unit: MICROSECOND
699+
`, nil)
700+
require.NoError(t, err)
701+
702+
encodeProc, err := newParquetEncodeProcessorFromConfig(encodeConf, nil)
703+
require.NoError(t, err)
704+
705+
commonSchema := &schema.Common{
706+
Type: schema.Object,
707+
Children: []schema.Common{
708+
{Name: "id", Type: schema.Int64},
709+
{Name: "ts", Type: schema.Timestamp},
710+
},
711+
}
712+
713+
msg := service.NewMessage([]byte(`{"id":1,"ts":"2026-04-17T12:00:00Z"}`))
714+
msg.MetaSetMut("benthos_schema", commonSchema.ToAny())
715+
716+
batches, err := encodeProc.ProcessBatch(t.Context(), service.MessageBatch{msg})
717+
require.NoError(t, err)
718+
require.Len(t, batches, 1)
719+
require.Len(t, batches[0], 1)
720+
721+
pqBytes, err := batches[0][0].AsBytes()
722+
require.NoError(t, err)
723+
724+
pqFile, err := parquet.OpenFile(bytes.NewReader(pqBytes), int64(len(pqBytes)))
725+
require.NoError(t, err)
726+
assert.Contains(t, pqFile.Schema().String(), "unit=MICROS")
727+
}
728+
729+
func TestParquetEncodeTimestampUnitInvalid(t *testing.T) {
730+
env := service.NewEnvironment()
731+
err := env.NewStreamBuilder().AddProcessorYAML(`
732+
parquet_encode:
733+
schema:
734+
- { name: id, type: INT64 }
735+
default_timestamp_unit: PICOSECOND
736+
`)
737+
require.Error(t, err)
738+
}
739+
650740
func TestParquetEncodeProcessorConfigLinting(t *testing.T) {
651741
configTests := []struct {
652742
name string

0 commit comments

Comments
 (0)