diff --git a/plugins/outputs/parquet/README.md b/plugins/outputs/parquet/README.md index 19627cd5da2c4..b5e190f00369e 100644 --- a/plugins/outputs/parquet/README.md +++ b/plugins/outputs/parquet/README.md @@ -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. diff --git a/plugins/outputs/parquet/parquet.go b/plugins/outputs/parquet/parquet.go index a4cf876d0c45b..1f230b5acd8ff 100644 --- a/plugins/outputs/parquet/parquet.go +++ b/plugins/outputs/parquet/parquet.go @@ -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" ) @@ -38,6 +39,7 @@ type Parquet struct { TimestampFieldName string `toml:"timestamp_field_name"` Log telegraf.Logger `toml:"-"` + root *os.Root metricGroups map[string]*metricGroup } @@ -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 @@ -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") } @@ -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) } @@ -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) } diff --git a/plugins/outputs/parquet/parquet_test.go b/plugins/outputs/parquet/parquet_test.go index b063ce0997605..f26993fa92dd7 100644 --- a/plugins/outputs/parquet/parquet_test.go +++ b/plugins/outputs/parquet/parquet_test.go @@ -3,6 +3,8 @@ package parquet import ( "os" "path/filepath" + "slices" + "strings" "testing" "time" @@ -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" ) @@ -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) +}