SeaTunnel Edge Agent 架构解析:WAL 出站队列、调度主循环与 EdgeSocket 协议边界设计
2026/9/17 20:17:32 网站建设 项目流程

SeaTunnel Edge Agent 架构解析:WAL 出站队列、调度主循环与 EdgeSocket 协议边界设计

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

SeaTunnel Edge Agent 是 SeaTunnel 中部署在边缘主机上的轻量采集器:它读取本地文件、通过 WAL(Write-Ahead Log)持久化出站队列缓冲数据,再以 EdgeSocket 行协议把批次推送到运行中的 Zeta 作业。本文基于仓库中的架构文档与对应源码,完整拆解其设计目标、逻辑架构、运行时行为、出站队列状态机、与 Engine 的协议契约以及可靠性语义,帮助你在部署边缘数据采集链路时理解“数据从哪里落盘、何时重放、何时可以认为送达”,并在排障时准确判断问题归属于 Agent 侧还是引擎侧。

1. 背景与设计目标

1.1 问题背景

在许多生产网络中,Zeta 集群无法直接访问边缘本地的文件(例如磁盘上的应用日志、NDJSON 路径)。此时需要一个专职的边缘采集进程来:

  • 在靠近数据源的位置读取本地记录;
  • 容忍网络间歇性中断;
  • 在不把引擎 worker 嵌入边缘站点的前提下,把数据投递给正在运行的 SeaTunnel 管道。

这正是 Edge Agent 的定位:它不是 SeaTunnel Engine 的替代品,而是“边缘侧接入”的专用进程。

1.2 设计目标

根据 架构文档,Edge Agent 的设计目标可以归纳为五点:

  1. 独立部署:打包与生命周期和引擎 worker 解耦,Agent 以独立进程运行在边缘主机上,安装布局见 Deployment Guide;
  2. 可持久化的出站缓冲:基于 WAL 的出站队列,具备显式的状态迁移(PENDING → SENDING → ACKED / DEAD);
  3. 与 Zeta 协议对齐:复用 EdgeSocket 行协议(__AUTH__/__BATCH__→ RECEIVED),引擎侧实现见 EdgeSocket source connector;
  4. 运维简单:YAML 配置 + 可预测的调度器循环;
  5. 边界清晰:发送路径与投递语义被显式限定,便于稳定运维与扩展。

1.3 架构定位:Edge Agent vs SeaTunnel Engine

维度Edge AgentSeaTunnel Engine
运行时位置边缘主机上的独立进程集群 Coordinator 与 worker 任务
输入访问边缘主机本地文件(glob 路径,含日志文件)连接器可达的系统(数据库、消息队列、对象存储等)
持久化与状态本地 WAL 出站队列与输入位置存储Checkpoint、作业状态、任务级恢复
网络角色主动拨号 Engine 的 EdgeSocket ingress 端点暴露 ingress 并编排内部任务执行
主要职责边缘数据接入与转发端到端数据集成与管道执行

1.4 职责边界:作为稳定契约

架构文档明确要求把以下边界当作稳定契约对待,这是理解整个系统设计的关键:

边界归属 Edge Agent归属 Engine / 作业
持久化边界持久化本地出站行,重试直至引擎返回 RECEIVED接入被接受之后的持久化处理
Checkpoint 边界本版本不感知 checkpoint(不使用__COMMIT__Checkpoint 生命周期与 exactly-once 语义
故障归属本地文件读取、本地 WAL 状态、传输层重连行为源端接入策略、下游 transform/sink 正确性

换言之,Agent 的合同在“引擎返回 RECEIVED”这一刻结束;再往后的正确性(去重、事务、checkpoint 语义)完全由作业负责。

2. 逻辑架构

2.1 部署拓扑

┌─────────────────────────────────────────────────────────────────┐ │ Edge host │ │ ┌───────────────────────────────────────────────────────────┐ │ │ │ Edge Agent process │ │ │ │ Input collector ──► Scheduler & batching │ │ │ │ │ │ │ │ │ │ │ ├──► Outbound queue (WAL) │ │ │ │ └──► Input position store (local persistence) │ │ │ │ │ │ │ │ │ └──► Transport client │ │ │ └───────────────────────────────────────────────────────────┘ │ └─────────────────────────────────────────────────────────────────┘ │ TCP line protocol ▼ ┌─────────────────────────────────────────────────────────────────┐ │ SeaTunnel Engine (Zeta) │ │ EdgeSocket Source ingress at configured output.endpoint │ └─────────────────────────────────────────────────────────────────┘

信任边界:Agent 信任本地文件系统访问及其自身的本地持久化存储;引擎信任__AUTH__中携带的配置好的 output.token。系统没有自动的集群端点发现——output.endpoint是静态配置的。

2.2 数据平面

记录在 Agent 内部通过单线程调度器循环流转:

控制平面(配置加载、插件选择、生命周期启停)只在启动和关闭时各执行一次;上述热路径即数据平面。

2.3 仓库模块划分

实现拆分为三个 Maven 模块(位于seatunnel-edge-agent/目录下):

层次职责模块
运行时核心YAML 解析、进程生命周期、调度器循环、出站队列与输入位置持久化seatunnel-edge-agent-starter
传输与编码EdgeSocket 或 console 输出、重连策略、RAW/PACKET 载荷模式seatunnel-edge-agent-transport
输入插件文件采集器、NDJSON 归一化、多行日志组装seatunnel-edge-agent-connector

从源码结构看,各模块边界与文档描述一一对应:

  • 调度器主循环与 WAL 实现位于 EdgeAgentRuntimeScheduler 与wal/包(含sqlite/mem/两套实现);
  • 行协议常量与重连逻辑位于 EdgeSocketProtocol 与socket/EdgeTransportClient
  • 文件采集、glob 解析、多行组装位于 FileCollectReader。

发行包由seatunnel-dist的 edge-agent 组装产出,下载方式见 Download。

3. 运行时行为

3.1 启动与关闭

启动时序(与 EdgeAgentRuntimeBootstrap 的实现一致):

从源码看有两个值得注意的实现细节:

  1. transport 先于 reader 打开EdgeAgentRuntimeBootstrap.start()中先调用ctx.getTransport().open()再调用ctx.getReader().open(),确保 ingress 就绪后才开始拉取数据;任何一步失败都会触发close()释放部分资源(见 EdgeAgentRuntimeBootstrap.java#L47-L61)。
  2. 优雅关闭时排空内存批runUntilStoppedfinally块中调用flushBufferToWal(),退出前把仍在 RAM 中的缓冲写入出站队列,避免记录丢失(见 EdgeAgentRuntimeScheduler.java#L96-L116)。

3.2 调度器主循环

每一轮调度器循环做五件事:

  1. 轮询输入——从配置的输入采集器(input.*)读取最多queue.poll-batch-size条事件;
  2. 内存缓冲——累积事件直到达到agent.bulk-max-sizeagent.flush-interval-ms超时;
  3. 刷入持久层——将事件以 PENDING 行追加到出站队列;同时持久化每个事件的输入位置(文件偏移量 / 行元数据);
  4. 发送出站——认领 PENDING 行、编码载荷、经 transport 发送;收到 RECEIVED 后标记为 ACKED;
  5. 事务性维护——将耗尽的 PENDING 行标记为 DEAD、复活(resurrect)滞留的 SENDING 行、按queue.acked-retention-ms删除已确认行、空闲时休眠agent.idle-sleep-ms

源码层面,单轮迭代由runOnce完成:先reader.poll(maxPollRecords),再按需flushBufferToWal(),最后sendClaimedRecords()(见 EdgeAgentRuntimeScheduler.java#L154-L206)。发送失败时行为与文档一致:

  • 一般 I/O 失败:行保持 SENDING,等待 resurrect/重试,同时按retry.backoff-ms~retry.backoff-max-ms指数退避后break本轮发送循环;
  • DECRYPT_FAILED 是致命错误:源码直接throw中断调度循环,对应文档中“fatal configuration error”的语义(见isDecryptFailed判断,EdgeAgentRuntimeScheduler.java#L208-L211)。

3.3 出站队列状态机

每条出站记录的完整状态定义在 WalRecordStatus 中(PENDING / SENDING / ACKED / DEAD),迁移规则如下:

PENDING │ claim for send (attempt_count++) ▼ SENDING │ 发送成功 + 引擎 RECEIVED ├──────────────► ACKED │ └ 发送失败 / 超时 / 在 RECEIVED 前崩溃 ▼ PENDING(attempt 计数加一;通过 resurrect 回到 PENDING) │ └ attempt_count >= retry.max-attempts ▼ DEAD(不再发送;由运维清理)
  • resurrect 机制resurrectSending周期性地把滞留的 SENDING 行拉回 PENDING,保证“认领之后、RECEIVED 之前崩溃”不会让数据卡死。主循环每轮检查是否到达queue.resurrect-interval-ms时间点后执行walStore.resurrectSending(...)
  • DEAD 机制markExceededAsDead把超过retry.max-attempts的行移入 DEAD,之后不再认领。

一个容易忽略但文档明确强调的点:文件位置是在事件追加进出站队列时(与 WAL 插入同一次 flush)持久化的,而不是在引擎确认批次时。从源码看,flushBufferToWal中每条事件先walStore.append(event)saveSourcePositionIfPresent(event)(见 EdgeAgentRuntimeScheduler.java#L221-L240)。恢复依赖 WAL 行 + 已保存的输入位置两者共同完成。

4. 出站队列与输入位置:双持久化存储

Agent 维护两个相互独立的本地持久化关注点:

存储用途更新时机重启后恢复
出站队列(WAL)在引擎对批次返回 RECEIVED 之前的持久性事件从内存刷入时;行在 PENDING → SENDING → ACKED 间迁移SENDING 行回退为 PENDING;未发送数据被重试
输入位置存储从断点继续读取本地文件,避免重读已持久化的事件与出站队列追加同一次 flush(按事件保存位置)采集器从保存的字节偏移量 / 行位置重新打开

把两者分开避免了“数据发到了网络上”与“磁盘上从哪里继续读”这两个问题的耦合:网络中断不应重置文件 tail 位置;反过来,推进文件游标也不意味着远端管道已经提交了数据

从源码结构看,这两个存储各有 SQLite 实现(SqliteWalStoreSqliteSourcePositionStore,位于seatunnel-edge-agent-starterwal/sqlite/包)与内存实现(MemWalStore/MemSourcePositionStore,用于测试和非持久模式),分别由WalStoreFactory/ SPI 选择。仓库中对应的单测如 SqliteWalStoreTest 与 SqliteSourcePositionStoreTest 覆盖了各自的持久化行为。

5. 网络与 EdgeSocket 契约

5.1 端点模型

transport 客户端连接的是静态配置output.endpoint(host:port),没有集群服务发现。变更接入主机意味着改配置并重启 Agent。

5.2 线上协议

Agent 实现了 EdgeSocket 的采集器一侧,本版本不发送__COMMIT__。Agent 侧的持久性以“收到 RECEIVED 后 WAL 行进入 ACKED”为终点,而不是轮询引擎 checkpoint。

协议常量定义在 EdgeSocketProtocol 中,与引擎侧文档逐一对应:

步骤Agent → EngineEngine → AgentAgent 处理方式
认证__AUTH__:<token>ACK / AUTH_FAILED / REJECTEDREJECTED:快速失败,不自动重连(说明存在重复采集器)
批次__BATCH__:<batchId>:<payload>RECEIVED / RETRY /QUEUE_FULL:<ms>/ DECRYPT_FAILEDQUEUE_FULL:等待并重发;DECRYPT_FAILED:致命的配置错误

ACK 与 RECEIVED 的区分:ACK 属于认证阶段;RECEIVED 属于批次接入阶段,且是唯一能把 WAL 行从 SENDING 推进到 ACKED 的成功响应。

batchId 的分配:线上__BATCH__:<batchId>:...中的 batchId 是 WAL 行的batch_id,从edge_agent_meta.next_batch_id单调分配(对应 SqliteBatchIdAllocator),它不是WAL 行主键id。源码中还有一层兼容处理:batchId = record.getBatchId() > 0 ? record.getBatchId() : record.getId()(见 EdgeAgentRuntimeScheduler.java#L180)。

批次发送时序:

5.3 重连策略

发生 I/O 失败(以及可重试的传输层失败)时,Agent 会:

  1. 使当前 socket 会话失效;
  2. 重试配置的端点候选(通常是一个静态主机);
  3. 重连、重新认证,并恢复认领待发送的出站行。

注意两套退避是分离的:transport 重连退避遵循output.*传输配置(如initial-backoff-ms/max-backoff-ms),而调度器重放节奏遵循retry.*。重连行为有专门测试覆盖:EdgeTransportClientReconnectTest。

5.4 协议契约与实现细节

稳定协议契约(跨版本承诺):

  • Agent 以__AUTH__:<token>认证,以__BATCH__:<batchId>:<payload>发送;
  • RECEIVED 表示该批次被接入接受;
  • 本版本 Agent 不发送__COMMIT__

当前运行时行为(实现细节,可能演进):

  • batchId 持久化为 WAL 行batch_id,从edge_agent_meta.next_batch_id分配;
  • 运行期使用单一调度器循环负责发送。

这一划分对运维很重要:写监控脚本、做二次开发时,应只依赖前者,避免把edge_agent_meta之类的内部实现当作接口。

6. 可靠性与投递语义

6.1 故障场景对照表

故障行为
Agent 进程崩溃重启时 SENDING 出站行恢复为 PENDING;内存批在未经过优雅关闭排空的情况下丢失
临时网络中断发送失败时行保持 SENDING;调度器退避;resurrectSending将滞留的 SENDING 行拉回 PENDING,然后重连并重试
采集端点变更更新output.endpoint并重启 Agent
优雅关闭退出前内存批刷入出站队列

6.2 投递模式

agent.delivery-guarantee未配置时默认为BEST_EFFORT(见 agent.yaml 中的注释:# delivery-guarantee: BEST_EFFORT # BEST_EFFORT: WAL + retry; NON: no WAL, stateless, drop on failure)。

BEST_EFFORT(默认)
  • Agent 维护本地 WAL 出站队列,重试直至引擎返回 RECEIVED(或在超过retry.max-attempts后行被标记 DEAD);
  • 同一 WAL 行在故障、认领与 RECEIVED 之间的崩溃、resurrectSending、或运维执行db wal-retry-dead可能被发送多次
  • Agent 不发送__COMMIT__;持久性止于引擎对批次返回 RECEIVED。

下游设计建议:把输出边界当作“可能重复投递”来处理——需要严格唯一性时使用幂等 sink 或去重键。完整参数见 Configuration — agent。

NON(无状态模式)
  • 没有 WAL、没有源位置持久化——Agent 完全无状态运行;
  • 事件从输入读取、内存批处理后直接经 transport 发送;发送失败时事件被丢弃并记录 warn 日志;
  • 重启后,文件读取依据输入配置(如read-from-beginning)重新开始,而不是保存的位置——之前发过的数据可能被重读重发;
  • queue.*retry.*配置节在 NON 模式下被忽略。

6.3 与 Engine Checkpoint 的边界

关注点Edge AgentSeaTunnel Engine
直到引擎 RECEIVED 的持久性WAL 支撑的出站队列
管道 exactly-once / checkpoint任务级 checkpoint 机制
向 Agent 回传提交游标不使用(不发送__COMMIT__

Agent 的合同在引擎对批次返回 RECEIVED 时结束;其后的正确性是作业的责任。

6.4 故障处置 Runbook

症状主要信号可能归属第一动作
AUTH_FAILEDAgent 传输/认证日志Agent + 作业配置对齐 output.token 与引擎 token,然后重启 Agent
REJECTEDAgent 传输/认证日志部署策略检查重复的采集器身份 / 监听策略冲突
积压增长(PENDING/SENDING)WAL 汇总与队列深度优先查 Agent 侧检查端点可达性、传输重试与引擎接入压力
反复出现 DEAD 行WAL 状态迁移Agent 配置 + 载荷兼容性检查死行、修复根因,再决定清除还是重试

7. 配置与扩展

7.1 配置面

运行时行为由单一agent.yaml驱动,顶层包含agentinputqueueretryoutput五个节。典型部署只需要配置input(带 paths)和生产output(transport + endpoint);queueretry不配置时采用默认值(sqlite-path: data/wal.db、内置重试策略)。

仓库与安装包中的完整示例文件位于 agent.yaml,安装根目录的config/下。关键默认值摘录:

# 身份与调度调优(可全部保持注释状态) agent: # id: auto # delivery-guarantee: BEST_EFFORT # BEST_EFFORT: WAL + 重试; NON: 无 WAL、无状态、失败即丢弃 # idle-sleep-ms: 200 # bulk-max-size: 256 # flush-interval-ms: 1000 # 采集对象:至少需要设置 paths input: paths: # 必填 — 待 tail 文件的 glob 模式 - "/var/log/*.log" # encoding: UTF-8 # read-from-beginning: false # glob-scan-interval-ms: 5000 # close-inactive-ms: 300000 # on-error: skip # 多行日志组装(取消注释启用): # multiline: # pattern: "^\\d{4}-\\d{2}-\\d{2}" # match: after # negate: false # max-lines: 500 # flush-idle-timeout-ms: 5000 # SQLite WAL 缓冲(delivery-guarantee 为 NON 时忽略) # queue: # sqlite-path: data/wal.db # poll-batch-size: 128 # cleanup-batch-size: 128 # acked-retention-ms: 0 # resurrect-batch-size: 100 # resurrect-interval-ms: 60000 # WAL 行发送重试策略(NON 模式忽略) # retry: # max-attempts: 16 # backoff-ms: 250 # backoff-max-ms: 300000 # 输出:默认 console(输出到 log/edge-agent.log,便于本地调试) output: type: console # --- 生产传输设置(取消注释启用) --- # type: transport # endpoint: "collector.example.com:9876" # type=transport 时必填 # token: "my-secret-token" # type=transport 时必填 # connect-timeout-ms: 5000 # read-timeout-ms: 30000 # initial-backoff-ms: 100 # max-backoff-ms: 30000 # max-reconnect-cycles: 16 # max-batch-send-attempts: 64 # packet-mode: RAW # compression: none # encryption: none

每个键的类型、默认值与校验规则,以权威的 Configuration Reference 为准。

7.2 输入

当前只实现了文件输入(input.type默认为file)。用input.paths配置 glob 即可 tail 本地文件:应用日志、NDJSON、轮转日志都表达为路径(例如/var/log/*.log),不存在独立的 log 或 event 输入类型。多行日志由 MultilineAssembler 按multiline.pattern组装。

参数与示例见 Input Configuration。新增的input.type取值是 SPI 扩展点,实现入口为EdgeInputReaderFactory(见 connector 模块 SPI 测试)。

7.3 输出与载荷编码

transport 与 console 的选择、与引擎端点的对齐、RAW 与 PACKET 两种载荷模式的适用场景,见 Output Configuration;线上响应的引擎侧定义见 EdgeSocket Source。

插件通过 YAML 中的input.typeoutput.type选择;新增一种输入或 transport 实现是扩展点行为,不改变调度器契约——调度器只面向EdgeInputReaderWalStoreEdgeCollectorTransport三个抽象编程,这正是单线程主循环能长期保持稳定的原因。

8. 小结

SeaTunnel Edge Agent 的架构可以浓缩为三句话:

  1. 一个单线程调度循环串联“读文件 → 内存批 → WAL PENDING → 认领发送 → RECEIVED → ACKED”,控制平面与数据平面彻底分离;
  2. 两个独立持久化存储——WAL 出站队列管“发没发出去”,输入位置存储管“从哪继续读”,互不耦合;
  3. 一条清晰的协议边界——Agent 只承诺到引擎 RECEIVED,__COMMIT__与 checkpoint 语义留给作业侧,投递语义默认 BEST_EFFORT(可重复投递),下游需按幂等/去重设计。

理解这三点后,部署调优(queue.*/retry.*/output.*参数)、故障定位(6.4 节 Runbook)和二次扩展(SPI 输入/输出插件)都有了确定的抓手。更多运维细节可继续阅读 Operations 与 FAQ。

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

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

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

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

立即咨询