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 的设计目标可以归纳为五点:
- 独立部署:打包与生命周期和引擎 worker 解耦,Agent 以独立进程运行在边缘主机上,安装布局见 Deployment Guide;
- 可持久化的出站缓冲:基于 WAL 的出站队列,具备显式的状态迁移(PENDING → SENDING → ACKED / DEAD);
- 与 Zeta 协议对齐:复用 EdgeSocket 行协议(
__AUTH__/__BATCH__→ RECEIVED),引擎侧实现见 EdgeSocket source connector; - 运维简单:YAML 配置 + 可预测的调度器循环;
- 边界清晰:发送路径与投递语义被显式限定,便于稳定运维与扩展。
1.3 架构定位:Edge Agent vs SeaTunnel Engine
| 维度 | Edge Agent | SeaTunnel 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 的实现一致):
从源码看有两个值得注意的实现细节:
- transport 先于 reader 打开:
EdgeAgentRuntimeBootstrap.start()中先调用ctx.getTransport().open()再调用ctx.getReader().open(),确保 ingress 就绪后才开始拉取数据;任何一步失败都会触发close()释放部分资源(见 EdgeAgentRuntimeBootstrap.java#L47-L61)。 - 优雅关闭时排空内存批:
runUntilStopped的finally块中调用flushBufferToWal(),退出前把仍在 RAM 中的缓冲写入出站队列,避免记录丢失(见 EdgeAgentRuntimeScheduler.java#L96-L116)。
3.2 调度器主循环
每一轮调度器循环做五件事:
- 轮询输入——从配置的输入采集器(
input.*)读取最多queue.poll-batch-size条事件; - 内存缓冲——累积事件直到达到
agent.bulk-max-size或agent.flush-interval-ms超时; - 刷入持久层——将事件以 PENDING 行追加到出站队列;同时持久化每个事件的输入位置(文件偏移量 / 行元数据);
- 发送出站——认领 PENDING 行、编码载荷、经 transport 发送;收到 RECEIVED 后标记为 ACKED;
- 事务性维护——将耗尽的 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 实现(SqliteWalStore与SqliteSourcePositionStore,位于seatunnel-edge-agent-starter的wal/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 → Engine | Engine → Agent | Agent 处理方式 |
|---|---|---|---|
| 认证 | __AUTH__:<token> | ACK / AUTH_FAILED / REJECTED | REJECTED:快速失败,不自动重连(说明存在重复采集器) |
| 批次 | __BATCH__:<batchId>:<payload> | RECEIVED / RETRY /QUEUE_FULL:<ms>/ DECRYPT_FAILED | QUEUE_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 会:
- 使当前 socket 会话失效;
- 重试配置的端点候选(通常是一个静态主机);
- 重连、重新认证,并恢复认领待发送的出站行。
注意两套退避是分离的: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 Agent | SeaTunnel Engine |
|---|---|---|
| 直到引擎 RECEIVED 的持久性 | WAL 支撑的出站队列 | — |
| 管道 exactly-once / checkpoint | — | 任务级 checkpoint 机制 |
| 向 Agent 回传提交游标 | 不使用(不发送__COMMIT__) | — |
Agent 的合同在引擎对批次返回 RECEIVED 时结束;其后的正确性是作业的责任。
6.4 故障处置 Runbook
| 症状 | 主要信号 | 可能归属 | 第一动作 |
|---|---|---|---|
| AUTH_FAILED | Agent 传输/认证日志 | Agent + 作业配置 | 对齐 output.token 与引擎 token,然后重启 Agent |
| REJECTED | Agent 传输/认证日志 | 部署策略 | 检查重复的采集器身份 / 监听策略冲突 |
| 积压增长(PENDING/SENDING) | WAL 汇总与队列深度 | 优先查 Agent 侧 | 检查端点可达性、传输重试与引擎接入压力 |
| 反复出现 DEAD 行 | WAL 状态迁移 | Agent 配置 + 载荷兼容性 | 检查死行、修复根因,再决定清除还是重试 |
7. 配置与扩展
7.1 配置面
运行时行为由单一agent.yaml驱动,顶层包含agent、input、queue、retry、output五个节。典型部署只需要配置input(带 paths)和生产output(transport + endpoint);queue与retry不配置时采用默认值(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.type与output.type选择;新增一种输入或 transport 实现是扩展点行为,不改变调度器契约——调度器只面向EdgeInputReader、WalStore、EdgeCollectorTransport三个抽象编程,这正是单线程主循环能长期保持稳定的原因。
8. 小结
SeaTunnel Edge Agent 的架构可以浓缩为三句话:
- 一个单线程调度循环串联“读文件 → 内存批 → WAL PENDING → 认领发送 → RECEIVED → ACKED”,控制平面与数据平面彻底分离;
- 两个独立持久化存储——WAL 出站队列管“发没发出去”,输入位置存储管“从哪继续读”,互不耦合;
- 一条清晰的协议边界——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),仅供参考