Telegraf Parquet 输出插件:把指标写入列式存储的完整实战指南
2026/9/14 4:00:20 网站建设 项目流程

Telegraf Parquet 输出插件:把指标写入列式存储的完整实战指南

【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf

Telegraf 的outputs.parquet插件负责将采集到的指标写入 Apache Parquet 列式存储文件,天然支持按指标名分组归档、时间轮转和 schema 自动推导,适合将监控指标沉淀到数据湖或供分析引擎直接查询。读完本文,你将掌握该插件的完整配置项、schema 生成与类型映射规则、文件轮转与关闭机制,以及如何用 Go/Python 工具快速探查生成的 Parquet 文件。

插件概述

outputs.parquet是 Telegraf 自 v1.32.0 起提供的一个datastore 类输出插件,适用于所有平台(💻 all)。它的工作方式非常简单直接:把收到的每条指标(metric)写入本地的.parquet文件,默认情况下按指标名分组,同名指标全部写入同一个文件。

[!IMPORTANT]如果一条指标的 schema 与文件中的 schema 不匹配,该指标将被丢弃。这是使用本插件时最重要的一条行为约束,后续的"Schema 生成与类型映射"一节会详细解释其成因与应对策略。

关于 Parquet 格式本身,官方文档(Apache Parquet、Parquet docs)提供了完整说明;如果你想了解用 Parquet 做毫秒级查询分析的实践,可参考 InfluxData 博客上关于 querying parquet 的文章。

注册与加载方式

与所有内置插件一样,该插件在 plugins/outputs/all/parquet.go 中被统一注册:

//go:build !custom || outputs || outputs.parquet import _ "github.com/influxdata/telegraf/plugins/outputs/parquet" // register plugin

它通过 parquet.go 中的outputs.Add("parquet", ...)注册为名为parquet的输出插件,并在初始化时把TimestampFieldName默认值设为"timestamp"

全局配置选项

所有 Telegraf 插件都支持一组全局与插件级配置项,用于修改指标、标签和字段、创建别名以及配置插件执行顺序。完整说明见 docs/CONFIGURATION.md(如aliasnamepass/namedroptagexcludemetricpass等),parquet 插件同样适用。

配置参数详解

完整可复制的示例配置在 plugins/outputs/parquet/sample.conf,其内容如下:

# A plugin that writes metrics to parquet files [[outputs.parquet]] ## Directory to write parquet files in. If a file already exists the output ## will attempt to continue using the existing file. # directory = "." ## Files are rotated after the time interval specified. When set to 0 no time ## based rotation is performed. # rotation_interval = "0h" ## Timestamp field name ## Field name to use to store the timestamp. If set to an empty string, then ## the timestamp is omitted. # timestamp_field_name = "timestamp"

三个配置项说明如下:

配置项类型默认值作用
directorystring"."写入 Parquet 文件的目录。若文件已存在,插件会尝试继续使用该文件
rotation_intervalduration"0h"基于时间的文件轮转间隔;设为0表示不进行时间轮转
timestamp_field_namestring"timestamp"存储指标时间戳的字段名;设为空字符串则省略时间戳列

directory:目录的处理与安全约束

从源码的Init()实现(parquet.go)可以看到目录的处理逻辑:

  • 未设置时默认使用当前工作目录"."
  • 会先转为绝对路径;
  • 目录不存在时以0750权限自动创建(os.MkdirAll),存在但不是目录则报错;
  • 目录会被os.OpenRoot打开为"根",所有文件创建/改名操作都限制在该目录内,无法通过../或绝对路径逃逸出配置目录。测试用例 TestCannotEscapeDirectory 专门验证了这一安全约束。

rotation_interval:时间轮转

  • 设置非 0 间隔后,每次Write都会检查当前文件,若文件修改时间 + 轮转间隔已超过当前时间,则先Close旧文件,再用相同文件名创建新 writer 继续写入(见源码rotateIfNeeded)。
  • 由于底层使用缓冲写入(buffered writer),不支持按大小轮转——文件可能不会在每个时间点真正落盘数据,无法可靠地以大小为依据切割。
  • 测试 TestRotation 通过设置RotationInterval: 1s验证了轮转行为。

timestamp_field_name:时间戳列

  • 插件为每个指标生成一行记录,并把m.Time().UnixNano()(纳秒整数,ArrowInt64类型)写入该列。
  • 设为空字符串""时,schema 中不包含时间戳列,测试 TestOmitTimestamp 验证了此时输出文件只有 1 列。
  • 同名冲突:如果某个字段或标签恰好也叫timestamp,在 schema 生成阶段会被丢弃并记录一条 Warning 日志(Ignoring the "timestamp" field or tag as that column holds the metric time)。解决方法是在配置中把timestamp_field_name改为其他名字,或重命名冲突的字段/标签。测试 TestTimestampFieldNameCollisionKeepsOneColumn 验证了冲突时只保留时间戳列且其物理类型为Int64

构建 Parquet 文件

Schema 生成

Parquet 文件在写入时必须携带 schema。Telegraf 的做法是:遍历同一分组内的全部指标,对所有字段和标签做并集,生成一个 Apache Arrow schema。具体规则如下(见源码createSchemagoToArrowType):

  1. 字段(field)优先:若某字段与某标签同名,字段胜出,其类型由字段值决定;测试 TestFieldTakesPrecedenceOverTagInAnyOrder 验证了无论 tag 在前还是 field 在前,都以字段类型为准。
  2. 标签一律映射为字符串(String)列
  3. 字段按 Go 类型映射为 Arrow 类型,支持的映射关系为:
Telegraf 字段类型(Go)Arrow 列类型
int8Int8
int16Int16
int32Int32
int64intInt64
uint8Uint8
uint16Uint16
uint32Uint32
uint64uintUint64
float32Float32
float64Float64
stringString
boolBoolean

其他类型(如数组、时间等)在goToArrowType中会返回unsupported type错误,导致对应指标被拒绝。测试 TestGather/data types 覆盖了全部 12 种支持类型。

  1. 所有列均可空(Nullable: true):某条指标缺少某列的值时,会写入 null 而不是 0 或空串。测试 TestMissingValuesReadBackAsNull 验证了这一点。

性能代价:由于 schema 需要额外的指标遍历,指标首次被刷写的那个 flush 周期会显著变慢;后续 flush 间隔会快得多。此外,schema 是在首次出现该指标名时确定的,若首次 flush 之后才出现新的字段,这些字段会被省略——这正是"schema 不匹配即被丢弃"规则的来源。

写入(Write)

写入链路使用buffered writerpqarrow.FileWriter.WriteBuffered,见源码Write方法)。它会把多次 flush 的指标先在内存中缓冲,再紧凑地合并写入同一个 Parquetrow group,从而获得更小的文件体积与更好的压缩效果。这种缓冲策略也直接决定了插件不支持基于文件大小的轮转。

关闭(Close)

Parquet 格式要求文件末尾带有footer(元数据脚注),因此文件必须被正确 Close,否则无法被正常读取。Close方法会遍历所有 metric group 的 writer 并逐一关闭,任何一个关闭失败都会汇总返回错误failed closing one or more parquet files(parquet.go)。

⚠️ 风险提示:如果 Telegraf 在写入 Parquet 文件的过程中崩溃,footer 可能尚未写出,该文件将损坏而无法读取。部署时建议为directory使用独立磁盘/持久化挂载,并配合监控及时发现异常退出。

文件命名

每个指标名对应的文件命名格式为(见源码Write):

{指标名}-{YYYY-MM-DD}-{unix秒时间戳}.parquet

例如cpu-2026-09-13-1789300000.parquet。指标名会直接拼入文件名,因此过长的指标名、含/\、tab、null 字节等非法文件名字符时,文件创建会失败,对应指标组整体被拒绝(测试 TestInvalidFilename 覆盖了这些边界情况)。

写入失败时的部分写语义

该插件实现了 Telegraf 的**部分写错误(PartialWriteError)**机制(定义见 internal/errors.go,规范见 docs/specs/tsd-008-partial-write-error-handling.md):

  • 单个指标值类型与 schema 列类型不匹配时(如 schema 为Int64却收到字符串),该条指标被reject,其余指标正常写入(见源码createRecordBatch中的类型校验,以及测试 TestPartialWrite/schema mismatch);
  • 文件创建/写入失败时,对应指标组的全部指标被 reject;
  • 被 reject 的指标会从输出缓冲中移除(不再重试),Write返回携带MetricsAccept/MetricsReject列表的PartialWriteError,供输出模型精确统计与跟踪。

文件轮转(File Rotation)

  • 启动时冲突处理:如果目标文件在启动时已存在,插件会先把现有文件重命名为带新时间戳的名字({指标名}-{新日期}-{新unix}.parquet),避免覆盖已有数据或产生 schema 冲突。逻辑见源码createWriter
  • 时间轮转:如配置rotation_interval,按上文所述基于文件修改时间进行轮转。
  • 不支持大小轮转:原因如前所述,buffered writer 使插件无法可靠感知文件实际数据量。

如何探查 Parquet 文件

方式一:Go CLI(arrow-go 官方工具)

Arrow 仓库提供了一个 Go 命令行工具,可快速读取与解析 Parquet 文件:

go install github.com/apache/arrow-go/v18/parquet/cmd/parquet_reader@latest parquet_reader <file>

方式二:Python pyarrow

也可以用 Python 的 pyarrow 快速打开并浏览文件:

import pyarrow.parquet as pq table = pq.read_table('example.parquet')

拿到 pyarrow.Table 后,即可调用其schemacolumn_namesto_pandas()num_rows等方法来进一步探查 schema 与数据内容。

典型配置示例

下面是一个真实可用的最小配置:把指标写入/var/lib/telegraf/parquet目录,每 6 小时轮转一次文件,并将时间戳列命名为ts以避免与业务字段冲突:

[[outputs.parquet]] directory = "/var/lib/telegraf/parquet" rotation_interval = "6h" timestamp_field_name = "ts"

需要注意的配合事项:

  • 若输入侧指标中包含名为tstimestamp的字段/标签,建议按上文规则重命名或调整timestamp_field_name,否则会被静默丢弃(仅记 Warning 日志);
  • 同一指标名的字段集合应尽量保持稳定,避免"先定 schema、后加字段"导致新字段被省略;
  • 输出目录需具备写入权限(插件会以0750自动创建),并考虑崩溃导致 footer 缺失的文件损坏风险。

总结

outputs.parquet插件为 Telegraf 提供了低门槛的列式存储落盘能力:按指标名分组归档、schema 自动推导、时间轮转与部分写错误处理一应俱全。上手时重点把握三点:首次写入即定型 schema(类型不匹配的指标会被丢弃)文件必须正常关闭才能读取(崩溃会损坏文件)时间戳列与业务字段同名会被忽略。结合 parquet.go 的源码与 parquet_test.go 的测试用例,你可以为数据湖/分析场景快速搭建一套可靠的指标归档管线。

【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询