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(如alias、namepass/namedrop、tagexclude、metricpass等),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"三个配置项说明如下:
| 配置项 | 类型 | 默认值 | 作用 |
|---|---|---|---|
directory | string | "." | 写入 Parquet 文件的目录。若文件已存在,插件会尝试继续使用该文件 |
rotation_interval | duration | "0h" | 基于时间的文件轮转间隔;设为0表示不进行时间轮转 |
timestamp_field_name | string | "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。具体规则如下(见源码createSchema与goToArrowType):
- 字段(field)优先:若某字段与某标签同名,字段胜出,其类型由字段值决定;测试 TestFieldTakesPrecedenceOverTagInAnyOrder 验证了无论 tag 在前还是 field 在前,都以字段类型为准。
- 标签一律映射为字符串(String)列。
- 字段按 Go 类型映射为 Arrow 类型,支持的映射关系为:
| Telegraf 字段类型(Go) | Arrow 列类型 |
|---|---|
int8 | Int8 |
int16 | Int16 |
int32 | Int32 |
int64、int | Int64 |
uint8 | Uint8 |
uint16 | Uint16 |
uint32 | Uint32 |
uint64、uint | Uint64 |
float32 | Float32 |
float64 | Float64 |
string | String |
bool | Boolean |
其他类型(如数组、时间等)在goToArrowType中会返回unsupported type错误,导致对应指标被拒绝。测试 TestGather/data types 覆盖了全部 12 种支持类型。
- 所有列均可空(Nullable: true):某条指标缺少某列的值时,会写入 null 而不是 0 或空串。测试 TestMissingValuesReadBackAsNull 验证了这一点。
性能代价:由于 schema 需要额外的指标遍历,指标首次被刷写的那个 flush 周期会显著变慢;后续 flush 间隔会快得多。此外,schema 是在首次出现该指标名时确定的,若首次 flush 之后才出现新的字段,这些字段会被省略——这正是"schema 不匹配即被丢弃"规则的来源。
写入(Write)
写入链路使用buffered writer(pqarrow.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 后,即可调用其schema、column_names、to_pandas()、num_rows等方法来进一步探查 schema 与数据内容。
典型配置示例
下面是一个真实可用的最小配置:把指标写入/var/lib/telegraf/parquet目录,每 6 小时轮转一次文件,并将时间戳列命名为ts以避免与业务字段冲突:
[[outputs.parquet]] directory = "/var/lib/telegraf/parquet" rotation_interval = "6h" timestamp_field_name = "ts"需要注意的配合事项:
- 若输入侧指标中包含名为
ts或timestamp的字段/标签,建议按上文规则重命名或调整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),仅供参考