编程 Agentic Streaming:七个运行时跑同一份 agent spec,四个 builder 方法是空操作

2026-09-29 00:04:09

Agentic Streaming:七个运行时跑同一份 agent spec,四个 builder 方法是空操作

项目信息

  • 仓库:https://github.com/Ugbot/Agentic-Streaming
  • 许可:Apache 2.0
  • 原名:Agentic Flink

Agentic Streaming 把 agent 当成流式、有状态、事件溯源的系统来构建:agent 的状态是有序事件日志上的物化视图,每个 conversation 只有一个 writer。Apache Flink 是功能最完整的运行时。同一份 agent spec 会在 Python、JVM、Clojure 的七个运行时上做一致性测试。

为什么用流处理引擎跑 Agent

真正干活的 agent(转账、处理工单、回答客户)要扛住流量峰值、节点故障和消息重放。流处理引擎已经解决了这些:带 key 的持久状态、exactly-once 或幂等处理、背压、自动恢复。把 agent 建在流处理引擎之上,意味着你原型的 agent 就是生产里跑的 agent,并且可以按规模挑引擎。

能构建什么

  • 基于实时事件流的 agent:Kafka、Postgres CDC、Redis pub/sub、webhook、Fluss、ZeroMQ 以及静态 seed 在 Flink 上都是 Channel\,多个 channel 可以扇入同一个 agent。NATS 是 Python 端口的后端,不是 Flink channel。
  • 带可校验结果的路由与链式:router -> path -> verifier 图,逐轮派发并校验回复,带输入/输出 guardrail,以及不依赖模型、可复现的规则大脑。
  • 几乎任何函数都能当工具:@Tool 方法、异步 ToolExecutor、MCP server(stdio 和 HTTP/SSE)、DJL 模型、HTTP 端点,都收进一个 ToolRegistry。
  • agent 调用 agent:A2A 把对端 agent 当工具,可以进程内、走 a2a-gateway(Agent Card、JSON-RPC、SSE;没有实现 gRPC 或 REST server),或作为显式的 pipeline 步骤,带重试和熔断。
  • 故障后仍存活的状态:per-conversation 内存加 keyed state,durability 由引擎提供,并有运行时专属测试证明:Flink checkpoint/savepoint、Pekko persistence、Clojure 上的 Datomic。
  • 引擎提供时就 exactly-once:Flink 的 checkpointed state;其余为幂等(effectively-once),以 ConversationStore 为真源。Kafka Streams adapter 没有配置 exactly_once_v2,这是 docs/portability/kafka-streams.md 里的设计说明,不是已发布代码。
  • saga 模式的长任务:后续步骤失败时,补偿处理器回滚多步流程;Pekko 运行时额外提供持久、可重试、human-in-the-loop 的工作流。
  • 跨事件模式检测(CEP):声明式 cep: 块(「5 分钟内同一主机 3 次异常就升级」)触发工具或派生事件;在每个 core 上都可移植、在 Flink 上原生,与定时器、窗口、重放、suspend/resume 并列。
  • 大多数数据系统:memory、vector、long-term storage 都是 SPI(Postgres、Redis/Valkey、Fluss、pgvector/Qdrant、NATS KV),通过 connection link 选择,交换时不碰 agent 代码。
  • 一份定义、多种部署:在 pipeline.yaml 里定义 agent,同一份 spec 跑在 Flink、Pekko、Clojure、Python core 或 experimental adapter 上。

运行时与一致性测试

契约是 spec/v1 下的 agentic/v1 spec,以及 spec/conformance/v1 下的 24 个 fixture。一个运行时只有用自己的 binding 跑完这些 fixture、并把结果报进生成的 capability matrix,才算通过一致性测试。共七个运行时和两个 facade binding:

  • reference: spec/tools/reference_runtime.py(fixture 的 oracle,不是生产运行时)
  • jvm-core: ports/jagentic-core JUnit ConformanceTest(不含 Flink 的 Java core)
  • flink: 根模块 JUnit FlinkConformanceTest(本地 MiniCluster 上的 Flink)
  • pekko: agentic-pekko JUnit PekkoConformanceTest(事件溯源 actor)
  • clojure: agentic-clj agentic.conformance/run-all(Datomic 上的纯 Clojure)
  • python: ports/pyagentic(纯 Python core)
  • pyflink: pyflink agentic_pyflink.conformance(PyFlink operator)
  • python-jvm、python-flink: python/agentic_flink(JPype facade,分别包 jvm-core 和 Flink)

每个能力通过、跳过还是失败,都在 docs/capabilities.md 里,由一次 fixture 运行重新生成。跳过的 fixture 是明确的「不支持」声明,永远不算通过。

崩溃持久性是另一回事。durable_store: supported 意味着需要它的三个 fixture(restart 后重放、suspend-resume、timer 在 restart 后存活)用 binding 声明的 store 通过了。跨真实进程或集群重启的持久性只由运行时专属测试证明(Flink WorkflowTurnFunctionMiniClusterTest;Pekko ConversationEntityTest/RedisJournalIT;Clojure datomic_test.clj;python test_agentic_runtime.py)。

experimental adapter,不做一致性测试

ports/experimental/ 下的 adapter(Faust、Kafka Streams、Temporal、Pulsar Functions、Ray、NATS JetStream、Quarkus、Spring、Celery、Dask、Airflow,带 gateway 的 Go core,以及 FastAPI gateway)早于 agentic/v1 spec。它们在自己的引擎上跑 banking 示例,共用 Python、Java 或 Go core。它们都不跑 agentic/v1 fixture,都不进 capability matrix,都不在验收路径上,有可能被移除。把它们当带可运行代码的设计研究,不是受支持的运行时。

快速开始

同一份 banking agent 可以跑在你喜欢的任意运行时上。一份 pipeline.yaml 描述带工具、知识库和 guardrail 的 router/path/verifier 图,到处都能原样运行。

git clone https://github.com/Ugbot/Agentic-Streaming.git && cd Agentic-Streaming

以下命令都在 fresh clone、Java 21 和 Python 3.11 或更高版本上跑过。JVM 线路用仓库里提交的 Maven wrapper(./mvnw);系统 mvn 低于 3.9 会被构建拒绝。

# Python,无模型、无基础设施
python -m pip install -e ports/pyagentic
PYTHONPATH=ports/agentic-pipeline python -m agentic_pipeline run examples/pipelines/banking.yaml --text "what is my balance?"

# 同一份 spec 跑 NATS JetStream adapter(需要 nats://127.0.0.1:4222 上的 NATS server;没有会以 ConnectionRefusedError 失败)
python -m pip install nats-py
PYTHONPATH=ports/agentic-pipeline:ports/experimental/nats python -m agentic_pipeline run examples/pipelines/banking.yaml --backend nats --text "card types?"

# Agentic Pekko:构建顺序有讲究,先装不含 Flink 的 Java core,再编译 Pekko
./mvnw -q -f ports/jagentic-core/pom.xml install -DskipTests
./mvnw -q -f agentic-pekko/pom.xml compile exec:java -Dexec.mainClass=org.jagentic.pekko.PipelineMain -Dexec.args="examples/pipelines/banking.yaml --text 'what is my balance?'"

# Agentic Clojure:Datomic 上的纯 Clojure(需要 Clojure CLI)
cd agentic-clj && clojure -M:run && cd ..

# Apache Flink:code-first 框架。编译约 400 个源文件,然后在本地 MiniCluster 上跑单元套件
./mvnw clean test

# 所有一等 JVM 模块按依赖顺序一次构建,含测试
./mvnw -f reactor/pom.xml verify

Python、Pekko、Clojure 三条线都返回 path payments 和余额 1234.56。

有两个旧版本列出的命令在 fresh clone 上跑不通,在代码 owner 修复前不属于快速开始(跟踪在 docs/audit-backlog.md,AGS-40):

  • ./mvnw exec:java -Dexec.mainClass=org.agentic.flink.example.QuickStartExample 在 AgentBuilder.build() 里以 "Initial state has no outgoing transitions" 失败,发生在调用任何模型之前。
  • ./mvnw exec:java -Dexec.mainClass=org.agentic.flink.pipeline.FlinkPipelineRunner 在 compile classpath 上失败,因为 flink-connector-datagen 是 test scope;在 test classpath 上提交 job 时以 "Could not deserialize stream node" 失败。
Agent agent = Agent.builder()
.withId("research-bot")
.withSystemPrompt("You are a research assistant.")
.withChatConnection(LangChain4jChatConnection.ollama("http://localhost:11434"))
.withChatSetup(ChatSetup.builder()
.withModel("qwen2.5:7b")
.withTemperature(0.3)
.withMaxResponseTokens(2048)
.withOutputSchema(OutputSchema.of(ResearchVerdict.class))
.build())
.withMcpServer(McpServerSpec.stdio("calc", "npx", "-y", "mcp-server-calculator"))
.withSkill(Skill.builder()
.withName("citations")
.withTools("doc-fetch", "summarize")
.withSystemPromptFragment("Prefer primary sources. Cite arxiv IDs.")
.build())
.withListener(new LoggingAgentEventListener(), new MetricsAgentEventListener())
.withMaxIterations(10)
.build();

每个 with* 方法都是可选的,默认值通过 ServiceLoader 发现。最小可用 agent 是 Agent.builder().withId(...).withSystemPrompt(...).build()。

有四个 builder 方法存在,但本仓库里没有任何 operator 读取:withShortTermTtl、withVectorMemory、withLongTermStore、withMemoryChannel。它们在 Agent 上存一个值,但没人消费,调用它们运行时不改变任何东西(见 docs/audit-backlog.md,AGS-32)。今天要用 Flink state 做短期记忆,就在你的 RichFunction.open() 里绑定 FlinkStateShortTermMemory.spec();向量记忆直接在 operator 里用 FlinkStateVectorMemory 或 FlinkStateHnswVectorMemory。

构建 agent:Python

from agentic_flink import load
spec = load("spec/conformance/v1/workflows/support.yaml")
result = spec.run(runtime="flink-jvm", text="what is my balance?", conversation_id="c1", turn_id="t1")

构建 agent:Agentic Pekko(actor)

agent 大脑从 Flink-free core 原样复用,只有 actor 和 persistence 外壳是 Pekko,每个 conversation 一个事件溯源、分片的实体。

mvn -f agentic-pekko/pom.xml exec:java -Dexec.mainClass=org.jagentic.pekko.http.HttpMain
curl -XPOST localhost:8080/agent -H 'content-type: application/json' -d '{"conversation_id":"c1","user_id":"u","text":"what is my balance?"}'

构建 agent:Agentic Clojure(Datomic)

地道的 Clojure 实现:brain、router、verifier 都是函数,transcript 是不可变的 Datomic datom,因此历史可以 time-travel。

(defn balance-brain [user-text ctx]
(str "[payments] Your balance is " (ctx/call-tool ctx "get_balance" {})))
clojure -M:run / -M:http / -M:time-travel

用 pipeline.yaml 换后端

Flink 是功能最丰富的运行时,但 agent 本身与引擎无关。先在嵌入式 local runtime 上原型,然后改一行 YAML 就能迁到流式、持久或批处理后端。

backend: nats
agent:
router:  { kind: keyword, default: general, rules: { payments: [balance], cards: [card] } }
paths:
payments: { brain: llm, prompt: "You are a payments specialist.", tools: [get_balance] }
cards:    { brain: rule, prompt: "You answer card questions." }
general:  { brain: rule, prompt: "You answer general questions." }
tools:   [ { id: get_balance, kind: constant, value: 1234.56 } ]
stores:  { conversation: { kind: redis, url: "${AGENTIC_REDIS_URL}" } }

Python loader 接受 backend: local | celery | nats,其它值抛 ValueError。外部服务(Redis/Valkey、Kafka/Fluss、Postgres、NATS)藏在接口之后,通过 examples/compose/externals.yml 起来。

模型:agent 是事件流上的物化视图

一个 conversation 不是 request/response 调用,而是一份有序事件日志(轮次、工具结果、模型输出、路由决策),agent 的状态就是重放这份日志得到的值。

  • Event sourcing:日志是唯一真源,状态是派生的。这就是 durability、replay、审计和恢复的来源,每个引擎实现方式不同(Flink checkpoint、Kafka/NATS offset、Pulsar BookKeeper、Pekko persistence、Temporal history)。
  • CQRS:命令(「处理这一轮」)是每个 conversation 单 writer 的有序变更,而查询(「当前答案或状态是什么」)是对视图的扇出读。分开以后,一个 conversation 既可以是持久的 keyed entity,也可以是流。

仓库组件

  • Flink 框架:Apache Flink 上的完整 agent 框架,包含 state-first memory、vector memory、CEP、chat/…

架构

核心是 spec/v1 的 agentic/v1 契约和 spec/conformance/v1 的 24 个 fixture。每个运行时用自己的 binding 跑 fixture,结果汇总进 capability matrix,作为跨运行时的能力对照面。状态模型统一是事件日志上的物化视图,每个 conversation 一个 writer;持久性由各引擎各自实现(Flink checkpoint、Kafka/NATS offset、Pulsar BookKeeper、Pekko persistence、Temporal history)。在幂等路径下,ConversationStore 是唯一真源。

推荐文章

程序员茄子在线接单