SeaTunnel Edge Agent 架构解析:WAL 出站队列、调度主循环与 EdgeSocket 协议边界设计 📅 发布时间:2026/9/17 20:17:51 👁 浏览次数: SeaTunnel Edge Agent 架构解析WAL 出站队列、调度主循环与 EdgeSocket 协议边界设计【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel Edge Agent 是 SeaTunnel 中部署在边缘主机上的轻量采集器它读取本地文件、通过 WALWrite-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 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 的实现一致从源码看有两个值得注意的实现细节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 前崩溃 ▼ PENDINGattempt 计数加一通过 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.endpointhost:port没有集群服务发现。变更接入主机意味着改配置并重启 Agent。5.2 线上协议Agent 实现了 EdgeSocket 的采集器一侧本版本不发送__COMMIT__。Agent 侧的持久性以“收到 RECEIVED 后 WAL 行进入 ACKED”为终点而不是轮询引擎 checkpoint。协议常量定义在 EdgeSocketProtocol 中与引擎侧文档逐一对应步骤Agent → EngineEngine → AgentAgent 处理方式认证__AUTH__:tokenACK / AUTH_FAILED / REJECTEDREJECTED快速失败不自动重连说明存在重复采集器批次__BATCH__:batchId:payloadRECEIVED / 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 会使当前 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 AgentSeaTunnel Engine直到引擎 RECEIVED 的持久性WAL 支撑的出站队列—管道 exactly-once / checkpoint—任务级 checkpoint 机制向 Agent 回传提交游标不使用不发送__COMMIT__—Agent 的合同在引擎对批次返回 RECEIVED 时结束其后的正确性是作业的责任。6.4 故障处置 Runbook症状主要信号可能归属第一动作AUTH_FAILEDAgent 传输/认证日志Agent 作业配置对齐 output.token 与引擎 token然后重启 AgentREJECTEDAgent 传输/认证日志部署策略检查重复的采集器身份 / 监听策略冲突积压增长PENDING/SENDINGWAL 汇总与队列深度优先查 Agent 侧检查端点可达性、传输重试与引擎接入压力反复出现 DEAD 行WAL 状态迁移Agent 配置 载荷兼容性检查死行、修复根因再决定清除还是重试7. 配置与扩展7.1 配置面运行时行为由单一agent.yaml驱动顶层包含agent、input、queue、retry、output五个节。典型部署只需要配置input带 paths和生产outputtransport endpointqueue与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 # typetransport 时必填 # token: my-secret-token # typetransport 时必填 # 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),仅供参考