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
5 changes: 5 additions & 0 deletions plugins/outputs/parquet/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,11 @@ of this occurring.

## File Rotation

Measurement names determine the file name and can be customized with
the [rename processor][rename].

[rename]: /plugins/processors/rename/README.md

If a file with the same target name exists at start, the existing file is
rotated to avoid over-writing it or conflicting schema.

Expand Down
122 changes: 82 additions & 40 deletions plugins/outputs/parquet/parquet.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import (

"github.com/influxdata/telegraf"
"github.com/influxdata/telegraf/config"
"github.com/influxdata/telegraf/internal"
"github.com/influxdata/telegraf/plugins/outputs"
)

Expand All @@ -38,6 +39,7 @@ type Parquet struct {
TimestampFieldName string `toml:"timestamp_field_name"`
Log telegraf.Logger `toml:"-"`

root *os.Root
metricGroups map[string]*metricGroup
}

Expand All @@ -50,15 +52,16 @@ func (p *Parquet) Init() error {
p.Directory = "."
}

stat, err := os.Stat(p.Directory)
if os.IsNotExist(err) {
if err := os.MkdirAll(p.Directory, 0750); err != nil {
return fmt.Errorf("failed to create directory %q: %w", p.Directory, err)
}
} else if !stat.IsDir() {
return fmt.Errorf("provided directory %q is not a directory", p.Directory)
if err := os.MkdirAll(p.Directory, 0750); err != nil {
return fmt.Errorf("failed to create directory %q: %w", p.Directory, err)
}

root, err := os.OpenRoot(p.Directory)
if err != nil {
return fmt.Errorf("failed to open directory %q: %w", p.Directory, err)
}

p.root = root
p.metricGroups = make(map[string]*metricGroup)

return nil
Expand All @@ -78,6 +81,11 @@ func (p *Parquet) Close() error {
}
}

if err := p.root.Close(); err != nil {
p.Log.Errorf("failed to close directory %q: %v", p.Directory, err)
errorOccurred = true
}

if errorOccurred {
return errors.New("failed closing one or more parquet files")
}
Expand All @@ -86,52 +94,86 @@ func (p *Parquet) Close() error {
}

func (p *Parquet) Write(metrics []telegraf.Metric) error {
groupedMetrics := make(map[string][]telegraf.Metric)
for _, metric := range metrics {
groupedMetrics[metric.Name()] = append(groupedMetrics[metric.Name()], metric)
grouped := make(map[string][]int)
for i, m := range metrics {
grouped[m.Name()] = append(grouped[m.Name()], i)
}

var writeErr internal.PartialWriteError

now := time.Now()
for name, metrics := range groupedMetrics {
if _, ok := p.metricGroups[name]; !ok {
filename := fmt.Sprintf("%s/%s-%s-%s.parquet", p.Directory, name, now.Format("2006-01-02"), strconv.FormatInt(now.Unix(), 10))
schema, err := p.createSchema(metrics)
if err != nil {
return fmt.Errorf("failed to create schema for file %q: %w", name, err)
}
writer, err := p.createWriter(name, filename, schema)
if err != nil {
return fmt.Errorf("failed to create writer for file %q: %w", name, err)
}
p.metricGroups[name] = &metricGroup{
builder: array.NewRecordBuilder(memory.DefaultAllocator, schema),
filename: filename,
schema: schema,
writer: writer,
}
for name, indices := range grouped {
batch := make([]telegraf.Metric, len(indices))
for j, i := range indices {
batch[j] = metrics[i]
}

group, err := p.groupFor(name, batch, now)
if err != nil {
writeErr.MetricsReject = append(writeErr.MetricsReject, indices...)
writeErr.MetricsRejectErrors = append(writeErr.MetricsRejectErrors, err)
continue
}

if p.RotationInterval != 0 {
if err := p.rotateIfNeeded(name); err != nil {
return fmt.Errorf("failed to rotate file %q: %w", p.metricGroups[name].filename, err)
return fmt.Errorf("failed to rotate file %q: %w", group.filename, err)
}
}

record, err := p.createRecordBatch(metrics, p.metricGroups[name].builder, p.metricGroups[name].schema)
record, err := p.createRecordBatch(batch, group.builder, group.schema)
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 create record for file %q: %w", group.filename, err)
}
err = group.writer.WriteBuffered(record)
record.Release()
if err != nil {
return fmt.Errorf("failed to write to file %q: %w", group.filename, err)
}
writeErr.MetricsAccept = append(writeErr.MetricsAccept, indices...)
}

return nil
if len(writeErr.MetricsReject) == 0 {
return nil
}
writeErr.Err = fmt.Errorf("rejected %d metric(s): %w", len(writeErr.MetricsReject), errors.Join(writeErr.MetricsRejectErrors...))

return &writeErr
}

func (p *Parquet) groupFor(name string, metrics []telegraf.Metric, now time.Time) (*metricGroup, error) {
if group, found := p.metricGroups[name]; found {
return group, nil
}

schema, err := p.createSchema(metrics)
if err != nil {
return nil, fmt.Errorf("failed to create schema for file %q: %w", name, err)
}

filename := parquetFilename(name, now)
writer, err := p.createWriter(name, filename, schema)
if err != nil {
return nil, err
}

group := &metricGroup{
builder: array.NewRecordBuilder(memory.DefaultAllocator, schema),
filename: filename,
schema: schema,
writer: writer,
}
p.metricGroups[name] = group

return group, nil
}

func parquetFilename(name string, at time.Time) string {
return fmt.Sprintf("%s-%s-%s.parquet", name, at.Format("2006-01-02"), strconv.FormatInt(at.Unix(), 10))
}

func (p *Parquet) rotateIfNeeded(name string) error {
fileInfo, err := os.Stat(p.metricGroups[name].filename)
fileInfo, err := p.root.Stat(p.metricGroups[name].filename)
if err != nil {
return fmt.Errorf("failed to stat file %q: %w", p.metricGroups[name].filename, err)
}
Expand Down Expand Up @@ -278,20 +320,20 @@ func (p *Parquet) createSchema(metrics []telegraf.Metric) (*arrow.Schema, error)
}

func (p *Parquet) createWriter(name, filename string, schema *arrow.Schema) (*pqarrow.FileWriter, error) {
if _, err := os.Stat(filename); err == nil {
now := time.Now()
rotatedFilename := fmt.Sprintf("%s/%s-%s-%s.parquet", p.Directory, name, now.Format("2006-01-02"), strconv.FormatInt(now.Unix(), 10))
if err := os.Rename(filename, rotatedFilename); err != nil {
if _, err := p.root.Stat(filename); err == nil {
if err := p.root.Rename(filename, parquetFilename(name, time.Now())); err != nil {
return nil, fmt.Errorf("failed to rename file %q: %w", filename, err)
}
}
file, err := os.Create(filename)

file, err := p.root.Create(filename)
if err != nil {
return nil, fmt.Errorf("failed to create file %q: %w", filename, err)
}

writer, err := pqarrow.NewFileWriter(schema, file, parquet.NewWriterProperties(), pqarrow.DefaultWriterProps())
if err != nil {
file.Close()
return nil, fmt.Errorf("failed to create parquet writer for file %q: %w", filename, err)
}

Expand Down
122 changes: 122 additions & 0 deletions plugins/outputs/parquet/parquet_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,8 @@ package parquet
import (
"os"
"path/filepath"
"slices"
"strings"
"testing"
"time"

Expand All @@ -11,6 +13,7 @@ import (

"github.com/influxdata/telegraf"
"github.com/influxdata/telegraf/config"
"github.com/influxdata/telegraf/internal"
"github.com/influxdata/telegraf/metric"
"github.com/influxdata/telegraf/testutil"
)
Expand Down Expand Up @@ -276,3 +279,122 @@ 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 TestUnusableMeasurementNamesAreRejected(t *testing.T) {
tests := []struct {
name string
rejected bool
}{
{"cpu", false},
{"disk.io", false},
{"..", false},
{"../evil", true},
{"../../../../tmp/evil", true},
{"/etc/evil", true},
{"", false},
{"nul\x00byte", true},
{strings.Repeat("a", 300), true},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
dir := t.TempDir()
p := &Parquet{Directory: dir, TimestampFieldName: "timestamp", Log: testutil.Logger{}}
require.NoError(t, p.Init())

err := p.Write([]telegraf.Metric{
metric.New(tt.name, nil, map[string]interface{}{"value": int64(1)}, time.Now()),
})
require.NoError(t, p.Close())

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

escaped, globErr := filepath.Glob(filepath.Join(filepath.Dir(dir), "*.parquet"))
require.NoError(t, globErr)
require.Empty(t, escaped)

if !tt.rejected {
require.NoError(t, err)
require.Len(t, written, 1)
return
}

var writeErr *internal.PartialWriteError
require.ErrorAs(t, err, &writeErr)
require.Equal(t, []int{0}, writeErr.MetricsReject)
require.Empty(t, writeErr.MetricsAccept)
require.Empty(t, written)
})
}
}

func TestSymlinkedMeasurementNameCannotEscapeDirectory(t *testing.T) {
dir := t.TempDir()
outside := t.TempDir()
require.NoError(t, os.Symlink(outside, filepath.Join(dir, "link")))

p := &Parquet{Directory: dir, TimestampFieldName: "timestamp", Log: testutil.Logger{}}
require.NoError(t, p.Init())

err := p.Write([]telegraf.Metric{
metric.New("link/escaped", nil, map[string]interface{}{"value": int64(1)}, time.Now()),
})
require.NoError(t, p.Close())

require.ErrorAs(t, err, new(*internal.PartialWriteError))

escaped, globErr := filepath.Glob(filepath.Join(outside, "*"))
require.NoError(t, globErr)
require.Empty(t, escaped)
}

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

require.ErrorAs(t, p.Write([]telegraf.Metric{
metric.New(strings.Repeat("a", 300), nil, map[string]interface{}{"value": int64(1)}, time.Now()),
}), new(*internal.PartialWriteError))

for i := 0; i < 8; i++ {
require.NoError(t, p.Write([]telegraf.Metric{
metric.New("good", nil, map[string]interface{}{"value": int64(i)}, time.Now()),
}))
}
require.NoError(t, p.Close())

written, err := filepath.Glob(filepath.Join(dir, "*.parquet"))
require.NoError(t, err)
require.Len(t, written, 1)
require.Contains(t, filepath.Base(written[0]), "good")
}

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

metrics := []telegraf.Metric{
metric.New("../evil", nil, map[string]interface{}{"value": int64(1)}, time.Now()),
metric.New("good", nil, map[string]interface{}{"value": int64(2)}, time.Now()),
metric.New("nul\x00byte", nil, map[string]interface{}{"value": int64(3)}, time.Now()),
metric.New("alsogood", nil, map[string]interface{}{"value": int64(4)}, time.Now()),
}

var writeErr *internal.PartialWriteError
require.ErrorAs(t, p.Write(metrics), &writeErr)
require.ElementsMatch(t, []int{0, 2}, writeErr.MetricsReject)
require.ElementsMatch(t, []int{1, 3}, writeErr.MetricsAccept)
require.Len(t, writeErr.MetricsRejectErrors, 2)
require.ElementsMatch(t,
[]int{0, 1, 2, 3},
append(slices.Clone(writeErr.MetricsAccept), writeErr.MetricsReject...),
)
require.NoError(t, p.Close())

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