Skip to content
Open
Show file tree
Hide file tree
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
3 changes: 3 additions & 0 deletions plugins/outputs/parquet/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,9 @@ When writing to a file, the schema is used to look for each value and if it is
not present a null value is added. The result is that if additional fields are
present after the first metric flush those fields are omitted.

Since column types are fixed at file creation, when an unknown value is logged
it will be logged as `null` and once per column.

### Write

The plugin makes use of the buffered writer. This may buffer some metrics into
Expand Down
169 changes: 88 additions & 81 deletions plugins/outputs/parquet/parquet.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ var defaultTimestampFieldName = "timestamp"

type metricGroup struct {
filename string
warned map[string]bool
builder *array.RecordBuilder
schema *arrow.Schema
writer *pqarrow.FileWriter
Expand Down Expand Up @@ -106,6 +107,7 @@ func (p *Parquet) Write(metrics []telegraf.Metric) error {
p.metricGroups[name] = &metricGroup{
builder: array.NewRecordBuilder(memory.DefaultAllocator, schema),
filename: filename,
warned: make(map[string]bool),
schema: schema,
writer: writer,
}
Expand All @@ -117,14 +119,13 @@ func (p *Parquet) Write(metrics []telegraf.Metric) error {
}
}

record, err := p.createRecordBatch(metrics, p.metricGroups[name].builder, p.metricGroups[name].schema)
group := p.metricGroups[name]
record := p.createRecordBatch(group, metrics)
err := group.writer.WriteBuffered(record)
record.Release()
if err != nil {
return fmt.Errorf("failed to create record for file %q: %w", p.metricGroups[name].filename, err)
}
if err = p.metricGroups[name].writer.WriteBuffered(record); err != nil {
return fmt.Errorf("failed to write to file %q: %w", p.metricGroups[name].filename, err)
return fmt.Errorf("failed to write to file %q: %w", group.filename, err)
}
record.Release()
}

return nil
Expand Down Expand Up @@ -154,89 +155,95 @@ func (p *Parquet) rotateIfNeeded(name string) error {
return nil
}

func (p *Parquet) createRecordBatch(metrics []telegraf.Metric, builder *array.RecordBuilder, schema *arrow.Schema) (arrow.RecordBatch, error) {
for index, col := range schema.Fields() {
func (p *Parquet) createRecordBatch(group *metricGroup, metrics []telegraf.Metric) arrow.RecordBatch {
for index, column := range group.schema.Fields() {
builder := group.builder.Field(index)

for _, m := range metrics {
if p.TimestampFieldName != "" && col.Name == p.TimestampFieldName {
builder.Field(index).(*array.Int64Builder).Append(m.Time().UnixNano())
continue
value := p.valueFor(m, column.Name)
if !appendValue(builder, value) {
p.warnOncef(
group, column.Name,
"Writing null for column %q of file %q as a %T value does not fit its %s column",
column.Name, group.filename, value, column.Type,
)
}
}
}

// Try to get the value from a field first, then from a tag.
var value any
var ok bool
value, ok = m.GetField(col.Name)
if !ok {
value, ok = m.GetTag(col.Name)
}
return group.builder.NewRecordBatch()
}

// if neither field nor tag exists, append a null value
if !ok {
switch col.Type {
case arrow.PrimitiveTypes.Int8:
builder.Field(index).(*array.Int8Builder).AppendNull()
case arrow.PrimitiveTypes.Int16:
builder.Field(index).(*array.Int16Builder).AppendNull()
case arrow.PrimitiveTypes.Int32:
builder.Field(index).(*array.Int32Builder).AppendNull()
case arrow.PrimitiveTypes.Int64:
builder.Field(index).(*array.Int64Builder).AppendNull()
case arrow.PrimitiveTypes.Uint8:
builder.Field(index).(*array.Uint8Builder).AppendNull()
case arrow.PrimitiveTypes.Uint16:
builder.Field(index).(*array.Uint16Builder).AppendNull()
case arrow.PrimitiveTypes.Uint32:
builder.Field(index).(*array.Uint32Builder).AppendNull()
case arrow.PrimitiveTypes.Uint64:
builder.Field(index).(*array.Uint64Builder).AppendNull()
case arrow.PrimitiveTypes.Float32:
builder.Field(index).(*array.Float32Builder).AppendNull()
case arrow.PrimitiveTypes.Float64:
builder.Field(index).(*array.Float64Builder).AppendNull()
case arrow.BinaryTypes.String:
builder.Field(index).(*array.StringBuilder).AppendNull()
case arrow.FixedWidthTypes.Boolean:
builder.Field(index).(*array.BooleanBuilder).AppendNull()
default:
return nil, fmt.Errorf("unsupported type: %T", value)
}
func (p *Parquet) valueFor(m telegraf.Metric, column string) interface{} {
if p.TimestampFieldName != "" && column == p.TimestampFieldName {
return m.Time().UnixNano()
}
if value, found := m.GetField(column); found {
return value
}
if value, found := m.GetTag(column); found {
return value
}

continue
}
return nil
}

switch col.Type {
case arrow.PrimitiveTypes.Int8:
builder.Field(index).(*array.Int8Builder).Append(value.(int8))
case arrow.PrimitiveTypes.Int16:
builder.Field(index).(*array.Int16Builder).Append(value.(int16))
case arrow.PrimitiveTypes.Int32:
builder.Field(index).(*array.Int32Builder).Append(value.(int32))
case arrow.PrimitiveTypes.Int64:
builder.Field(index).(*array.Int64Builder).Append(value.(int64))
case arrow.PrimitiveTypes.Uint8:
builder.Field(index).(*array.Uint8Builder).Append(value.(uint8))
case arrow.PrimitiveTypes.Uint16:
builder.Field(index).(*array.Uint16Builder).Append(value.(uint16))
case arrow.PrimitiveTypes.Uint32:
builder.Field(index).(*array.Uint32Builder).Append(value.(uint32))
case arrow.PrimitiveTypes.Uint64:
builder.Field(index).(*array.Uint64Builder).Append(value.(uint64))
case arrow.PrimitiveTypes.Float32:
builder.Field(index).(*array.Float32Builder).Append(value.(float32))
case arrow.PrimitiveTypes.Float64:
builder.Field(index).(*array.Float64Builder).Append(value.(float64))
case arrow.BinaryTypes.String:
builder.Field(index).(*array.StringBuilder).Append(value.(string))
case arrow.FixedWidthTypes.Boolean:
builder.Field(index).(*array.BooleanBuilder).Append(value.(bool))
default:
return nil, fmt.Errorf("unsupported type: %T", value)
}
}
func (p *Parquet) warnOncef(group *metricGroup, column, format string, args ...interface{}) {
if group.warned[column] {
return
}
group.warned[column] = true
p.Log.Warnf(format, args...)
}

func appendValue(builder array.Builder, value interface{}) bool {
switch v := value.(type) {
case nil:
builder.AppendNull()
return true
case int8:
return appendTyped(builder, v)
case int16:
return appendTyped(builder, v)
case int32:
return appendTyped(builder, v)
case int64:
return appendTyped(builder, v)
case int:
return appendTyped(builder, int64(v))
case uint8:
return appendTyped(builder, v)
case uint16:
return appendTyped(builder, v)
case uint32:
return appendTyped(builder, v)
case uint64:
return appendTyped(builder, v)
case uint:
return appendTyped(builder, uint64(v))
case float32:
return appendTyped(builder, v)
case float64:
return appendTyped(builder, v)
case string:
return appendTyped(builder, v)
case bool:
return appendTyped(builder, v)
default:
builder.AppendNull()
return false
}
}

func appendTyped[T any](builder array.Builder, value T) bool {
column, ok := builder.(interface{ Append(T) })
if !ok {
builder.AppendNull()
return false
}
column.Append(value)

record := builder.NewRecordBatch()
return record, nil
return true
}

func (p *Parquet) createSchema(metrics []telegraf.Metric) (*arrow.Schema, error) {
Expand Down
58 changes: 58 additions & 0 deletions plugins/outputs/parquet/parquet_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,9 @@ import (
"testing"
"time"

"github.com/apache/arrow-go/v18/arrow"
"github.com/apache/arrow-go/v18/arrow/array"
"github.com/apache/arrow-go/v18/arrow/memory"
"github.com/apache/arrow-go/v18/parquet/file"
"github.com/stretchr/testify/require"

Expand Down Expand Up @@ -276,3 +279,58 @@ func TestMissingValuesReadBackAsNull(t *testing.T) {
require.Equalf(t, int16(1), column.MaxDefinitionLevel(), "column %q is required, nulls would be written as zero", column.Name())
}
}

func TestConflictingValueTypesDoNotPanic(t *testing.T) {
dir := t.TempDir()
p := &Parquet{Directory: dir, TimestampFieldName: "timestamp", Log: testutil.Logger{}}
require.NoError(t, p.Init())

require.NoError(t, p.Write([]telegraf.Metric{
metric.New("demo", nil, map[string]interface{}{"k": int64(1)}, time.Now()),
}))
require.NoError(t, p.Write([]telegraf.Metric{
metric.New("demo", map[string]string{"k": "v"}, map[string]interface{}{"other": int64(2)}, time.Now()),
}))
require.NoError(t, p.Close())

written, err := filepath.Glob(filepath.Join(dir, "*.parquet"))
require.NoError(t, err)
require.Len(t, written, 1)

reader, err := file.OpenParquetFile(written[0], false)
require.NoError(t, err)
defer reader.Close()
require.Equal(t, int64(2), reader.NumRows())
}

func TestEveryConvertibleTypeRoundTrips(t *testing.T) {
values := map[string]interface{}{
"int8": int8(1), "int16": int16(2), "int32": int32(3), "int64": int64(4), "int": 5,
"uint8": uint8(6), "uint16": uint16(7), "uint32": uint32(8), "uint64": uint64(9), "uint": uint(10),
"float32": float32(11), "float64": float64(12),
"string": "thirteen", "bool": true,
}

for name, value := range values {
t.Run(name, func(t *testing.T) {
datatype, err := goToArrowType(value)
require.NoError(t, err)

builder := array.NewBuilder(memory.DefaultAllocator, datatype)
defer builder.Release()

require.True(t, appendValue(builder, value))
require.Equal(t, 0, builder.NullN())
})
}
}

func TestAppendValueNullsWhatItCannotWrite(t *testing.T) {
builder := array.NewBuilder(memory.DefaultAllocator, arrow.PrimitiveTypes.Int64)
defer builder.Release()

require.False(t, appendValue(builder, "not an int"))
require.False(t, appendValue(builder, []string{"unsupported"}))
require.True(t, appendValue(builder, nil))
require.Equal(t, 3, builder.NullN())
}
Loading