采集配置Pipeline演化路线 #1928
Takuka0311
started this conversation in
Ideas
Replies: 0 comments
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
-
整体演进顺序为 A → B → C,支线 D 与 B/C 交叉推进。
总览:三条主线 + 一条支线
StructureType: v2可安全运行PipelineEventGroup;删 v1 Runner;Go→C++ 用 PB 替代SendPb(LogGroup)分阶段依赖总览
flowchart TB subgraph P1["阶段 1 · 短期"] A12["A1–A2\n扩 Input 白名单"] B12["B1–B2\n透传契约 + 插件矩阵盘点"] end subgraph P2["阶段 2 · 中期"] B34["B3–B4\nFlusher Export / Processor Process 迁移\nflusher_sls v2 包装或双接口"] D12["D1–D3\n自监控 PB 推送 + 指标目录对齐"] A3["A3\n移除 Input 白名单"] A4["A4\n多 Input E2E"] B56["B5–B6\n插件 lint + 混合 Log/Metric/Span E2E"] end subgraph P3["阶段 3 · 长期"] C123["C1–C3\nSendPipelineEventGroup\nflusher_sls Export + C++ 接 PB"] D45["D4–D6\n废弃 GetGoMetrics/Alarms CGO\nWindows 全量采集"] C4["C4\n删 pluginv1Runner / Legacy 转换"] C5["C5\n默认 StructureType: v2"] C6["C6\n清理 LogGroup 热路径"] end subgraph Target["终态"] T["全部 Native Input → Go\n仅 v2 Runner + 全插件多事件\n双向 PipelineEventGroup PB"] end A12 --> A3 A12 --> A4 B12 --> B34 B12 --> D12 B34 --> B56 B34 --> C123 D12 --> D45 B34 --> D45 D45 --> C4 A3 --> T A4 --> T B56 --> C123 C123 --> C4 B34 --> C4 C4 --> C5 C5 --> C6 C6 --> T D45 --> TStructureType主线 A:所有 C++ Input 均可进入 Go Pipeline
现状(代码锚点)
core/config/CollectionConfig.cpp中hasFileInput/hasSecurityInput/hasSelfMonitorInput决定是否允许「Native Input + Go Processor/Flusher」。未命中白名单则报错extended * plugins coexist with native input plugins other than ...(约 L365、L495)。IsFlushingThroughGoPipeline()==false的 Pipeline(如input_internal_metrics→flusher_file)。目标
任意注册在 C++
PluginRegistry的 Continuous Native Input(input_forward、input_prometheus、input_container_stdio等),只要 Pipeline 含 Go Processor 或 Go Flusher,处理完成后均可进入 Go,不再维护 Input 类型白名单。分阶段实施
CollectionConfig.cpp;UT:CollectionConfigUnittest增例core/plugin/input/*;ProtocolConversion单测hasFileInput等分支及input_file仅允许单 Input 等过渡约束(代码内 TODO)CollectionConfig.cpp、CollectionPipeline.cpp(DefaultLogQueueSize特例)input_forward/input_prometheus+flusher_prometheus(v2)test/e2e/test_cases/*flowchart LR subgraph CppIn["C++ Native Input(目标:全部)"] IF["input_file"] IP["input_prometheus"] II["input_internal_metrics"] IF2["input_forward / ..."] end subgraph Gate["配置层(逐步拆除)"] CC["CollectionConfig.Parse\n白名单 → 默认允许"] end subgraph Bridge["数据层(已具备)"] PR["ProcessorRunner\nTransferPipelineEventGroupToPB"] PEG["ProcessPipelineEventGroup"] end CppIn --> CC --> PR --> PEG风险与约束: 多 Input 同 Pipeline、Native+Go 多 Flusher 共存规则(
nativeFlusherCnt、flusher_sls)需在 A3 一并梳理,避免仅拆白名单导致配置组合爆炸。主线 B:Processor / Flusher 支持全部 PipelineGroupEvents
现状(接口与实现差距)
fetchPluginVersion在未写global.StructureType时默认v1→pluginv1Runner(logstore_config.go)。pkg/pipeline/processor.go/flusher.go):ProcessLogs/Flush(LogGroup)— 仅 Log。Process(PipelineGroupEvents)/Export([]*PipelineGroupEvents)— Log / Metric / Span。plugins/processor/*大量仅ProcessLogs(约 50+ 处);v2 需Process,且多数假定models.Log。flusher_prometheus等仅 v2Export;flusher_http双接口;flusher_sls仅FlusherV1.Flush,pluginv2Runner.addFlusher无法直接加载。Record已对齐;v1 路径仍聚合为LogGroup供 slsFlush。目标
任一 v2 Processor/Flusher 对
PipelineGroupEvents:支持(按语义处理子集)、透传(不静默丢弃)、声明SupportedEventKinds(或等价机制)。分阶段实施
PassThroughEvents)pkg/pipeline/pkg/helperProcess/Export/ProcessLogsFlusherV2.Export;Log-only 对 Metric/Span 透传plugins/flusher/*ProcessLogs→Process(Log),Metric/Span 透传;flusher_slsv2 包装或双接口(过渡期仍可用SendPb)plugins/processor/*、plugins/flusher/slsProcess仅按Events[0]假定 Logflusher_prometheus/flusher_httptest/e2e/*flowchart TB subgraph Input["PipelineGroupEvents"] L["Log"] M["Metric"] S["Span"] end subgraph Proc["Processor 链(目标)"] P1["Processor A\n处理 Log"] P2["Processor B\n处理 Metric"] P3["Processor C\n透传 Span"] end subgraph Flush["Flusher(目标)"] FP["flusher_prometheus\nMetric"] FH["flusher_http\nLog / 告警"] end L --> P1 --> FH M --> P2 --> FP S --> P3 --> FHProcessLogsProcess只改 Log;Metric/Span 透传processor_log_to_sls_metric)FlushExport+ Logflusher_slsFlush→SendPbSendPipelineEventGroup主线 C:彻底废弃 LogGroup,统一 Pipeline 事件模型
现状:两套模型并存
v1StructureType: v2ProcessPipelineEventGroupReceivePipelineEventGroupTransferPBToLogGroupForV1(Metric/Span 丢弃)TransferPBToPipelineGroupEventsflusher_sls:Flush→SendPbExport,不回 C++pipeline->Send()→FlusherSLS::Sendflusher_sls与 Export 型 Flusherflusher_slsFlusherSLS::Send)flusher_slsGenerateGoPlugin注入 Go sls)LogGroup→SendPb→ C++)flusher_prometheus(v2)目标
SendPipelineEventGroup(PB),对称于ProcessPipelineEventGroup。pluginv1Runner、TransferPBToLogGroupForV1、ReceiveLogGroup等。StructureType: v2(在 B4 完成后)。sls_logs.LogGroup热路径。分阶段实施
SendPipelineEventGroup;与SendPb并存plugin_export.go、LogtailPlugin.cppflusher_sls实现FlusherV2.Export→SendPipelineEventGroupplugins/flusher/sls/flusher_sls.goFlusherSLS::Send接PipelineEventGroupPBcore/plugin/flusher/sls/*pluginv1Runner、Legacy PB 转换pluginmanager/*、pipeline_event_helper.goStructureType: v2;文档/示例切换logstore_config.goLogGroup热路径;协议保留兼容层pkg/protocol/sls_logs.*flowchart TB subgraph Today["Go → C++ 今日"] F1["flusher_sls.Flush"] LG["LogGroup.Marshal"] SP["SendPb / SendPbV2"] FS["FlusherSLS::Send"] F1 --> LG --> SP --> FS end subgraph Future["Go → C++ 目标(C1–C3)"] F2["flusher_sls.Export"] PB2["PipelineEventGroup PB"] SP2["SendPipelineEventGroup"] FS2["FlusherSLS 接 PB"] F2 --> PB2 --> SP2 --> FS2 end Today -.-> Future与 #2568 相关配置在 C5 完成前 须继续 显式
StructureType: v2。支线 D:统一 C++ 与 Go 自监控(指标 + 告警)
问题研判(与代码对照)
你提出的两点 成立,且代码里已有明确佐证;另有几项关联问题值得一并纳入支线范围。
(1)CGO 自定义结构体不稳定 —— 成立,Windows 上更严重
GetGoMetrics→ Gomalloc构造PluginMetrics/InnerPluginMetric/InnerKeyValue→ C++ 侧逐个free(LogtailPlugin.cppL476–505)GetGoAlarms→InnerGoAlarms/InnerGoAlarm同理(L508–575)→ 写入AlarmManager::SendExternalAlarmPipelineEventGroup不是同一路径ReadMetrics::UpdateMetrics()在#ifndef _MSC_VER下才调用GetGoMetrics;注释写明 「windows 上的 cgo 内存有问题…可能会 crash」(MetricManager.cppL168–184)业务数据 C++→Go 的
ProcessPipelineEventGroup(PipelineEventGroupPB)已打通;自监控 反向 仍走上述 CGO 结构体,并未统一。(2)Go 与 C++ 监控能力未对齐 —— 成立,有多层表现
WriteMetrics快照 →MetricsRecord链表 →SelfMonitorMetricEventGetGoMetrics→map内嵌 JSON 字符串(labels/counters/gauges)再解析METRIC_CATEGORY_UNKNOWNagent/runner/pipeline/plugin/plugin_source/componentmetric_category=plugin(GetPluginCommonLabels)plugin_source;与 C++ Input 级指标(文件偏移、采集延迟等)无法同维度对比ExportMetricRecords仅导出 Counter + Gauge(metrics_record.goL68–73)MetricsRecordRegisterMetricRecord;未注册则无数据AlarmManager::FlushAllRegionAlarm→PipelineEventGroup→ 自监控 PipelineAlarmManager,再与 C++ 告警合并 flush支线目标
PipelineEventGroupPB(Metric 用多值UntypedMultiDoubleValues,告警用 LogEvent + metadata),废弃GetGoMetrics/GetGoAlarms及InnerPlugin*/InnerGoAlarm结构体。plugin/plugin_source/runner语义一致,Windows/Linux 行为一致。input_internal_*→ Go Flusher;采集改为 Go 主动 Push(或 C++ 拉 PB 字节流),与ProcessPipelineEventGroup对称。分阶段实施
MetricConstants↔ Gopkg/selfmonitor);明确pluginvsplugin_source映射规则AlarmManager字段);复用TransferPipelineEventGroupToPB/ GoTransferMetricEventToPBprotobuf_public、pkg/helperPushSelfMonitorToCore(pb []byte)(或定时调用 C++ 注册的ReceiveSelfMonitorEventGroup);C++ 解析 PB 后直接SelfMonitorServer::PushSelfMonitorMetricEvents/ 告警入队plugin_export.go、SelfMonitorServer.cpp、MetricManager.cppplugin_source级 MetricsRecord(Input 插件);补齐 Counter/Gauge 以外类型的导出或等价聚合plugin_wrapper*.go、pkg/selfmonitorGetGoMetrics/GetGoAlarmsCGO;移除MetricManager.cpp的_MSC_VER屏蔽;Windows 回归LogtailPlugin.*、plugin_export.goflusher_file/Prometheus 可查到;Windows CI 覆盖test/e2e、core/unittest/monitorflowchart TB subgraph Today["采集 · 今日(拉模式 + CGO 结构体)"] TMR["SelfMonitorServer::SendMetrics"] RM["ReadMetrics::UpdateMetrics"] GM["GetGoMetrics / GetGoAlarms\nInnerPlugin* / InnerGoAlarm"] GO1["Go: ExportMetricRecords → map+JSON"] PARSE["C++: SelfMonitorMetricEvent(map)\nAlarmManager::SendExternalAlarm"] TMR --> RM --> GM --> GO1 --> PARSE WIN["Windows: GetGoMetrics 整段禁用"] GM -.-> WIN end subgraph Target["采集 · 目标(推模式 + PB)"] TMR2["SelfMonitorServer::SendMetrics"] PUSH["Go PushSelfMonitor PB\n或 C++ Pull PB bytes"] PB["PipelineEventGroup\nMetrics + Logs"] MERGE["C++ 直接 Merge 到\nSelfMonitorMetricEventMap / 告警队列"] TMR2 --> MERGE PUSH --> PB --> MERGE end subgraph Export["导出 · 已打通"] EX["input_internal_* → PB → Go Flusher"] end Today -.->|"D2–D5"| Target MERGE --> Exportflowchart LR subgraph Align["D1/D4 · 能力对齐"] CAT["category 统一\nplugin / plugin_source / runner"] NAME["指标名常量统一"] TYPE["Go Export 覆盖 C++ 指标类型"] end subgraph Transport["D2/D3/D5 · 传输统一"] PB2["PipelineEventGroup PB"] NOP["移除 CGO 结构体"] end Align --> Transportinput_internal_*导出plugin_source改造SendPipelineEventGroup基础设施;D5 与 C4 同期删除 Legacy 拉取/转换里程碑 / Progress
快照更新于 2026-07-17T09:08:29.123Z(来源:Epic Console board)
Beta Was this translation helpful? Give feedback.
All reactions