1. 引言
本文基于 Akka Actor 模型,演示如何模拟一条完整的 Kafka 消息处理流水线:从消费 Kafka 消息开始,经过加工处理,最终发送到名为 flink_out 的 Kafka 主题。示例使用 JDK 1.8 + Akka 2.5.11 + Scala 2.11,通过两个 ActorSystem 和两条 LookupEventBus 实现消息的解耦流转,完整复刻了真实项目中 DatasetServiceV2 / Planner / OutputServiceV2 的核心架构。
2. 整体架构与数据流
本示例的数据流如下:
- KafkaConsumerActor:模拟从 Kafka 消费原始订单消息(≈ DatasetServiceV2.process)。
- OrderProcessBus:LookupEventBus,按业务线 ID 将消息投递给对应的 Planner 路由器。
- OrderPlanner:RoundRobinPool 路由器(2~5 个 routee 动态扩缩),每个 Planner 内部串联一条 Stepper 处理链。
- Stepper 链:ValidateStepper(校验拆分)→ PriceStepper(模拟远程取价)→ ResultStepper(组装结果)→ OutputStepper(发到输出总线)。
- OutputBus:LookupEventBus,按输出通道 ID 将结果投递给输出管道。
- OutputPipeline:RoundRobinPool 路由器,模拟 KafkaProducer 将消息发送到 flink_out 主题。
flowchart TD
A[KafkaConsumerActor] -->|publish| B[OrderProcessBus]
B -->|按业务线ID投递| C[OrderPlanner Router]
C --> D[ValidateStepper]
D --> E[PriceStepper]
E --> F[ResultStepper]
F --> G[OutputStepper]
G -->|publish| H[OutputBus]
H -->|按通道ID投递| I[OutputPipeline Router]
I -->|模拟KafkaProducer.send| J[flink_out Kafka主题]
3. 项目结构与依赖
示例代码位于 src/test 目录,走 Maven 测试类路径即可运行,无需真实 Kafka / DB。核心依赖如下:
<dependency>
<groupId>com.typesafe.akka</groupId>
<artifactId>akka-actor_2.11</artifactId>
<version>2.5.11</version>
</dependency>
<dependency>
<groupId>com.typesafe</groupId>
<artifactId>config</artifactId>
<version>1.3.2</version>
</dependency>
4. 配置文件 demo.conf
配置文件为两个演示 ActorSystem 各分配一段配置,通过 ConfigFactory.load("demo") 加载,每个 system 用 ConfigFactory 取自己的子配置段。段内 akka.actor.deployment."/*" 指定路由器类型与动态扩缩容:
# 订单处理流水线(对应 DatasetProcessSystem)
OrderDemoProcessSystem {
akka.actor.deployment {
"/*" {
router = round-robin-pool
optimal-size-exploring-resizer {
enabled = on
lower-bound = 2
upper-bound = 5
}
}
}
}
输出管道(对应 OutputSystem)
OrderDemoOutputSystem {
akka.actor.deployment {
"/*" {
router = round-robin-pool
optimal-size-exploring-resizer {
enabled = on
lower-bound = 1
upper-bound = 3
}
}
}
}
5. 消息定义
定义订单消息和输出消息,作为 Actor 之间传递的数据载体:
import java.io.Serializable;
/** 订单消息:模拟从 Kafka 消费到的原始数据 */
public class OrderMessage implements Serializable {
private final long orderId;
private final long businessId;
private final String productName;
private final int quantity;
private final double unitPrice;
public OrderMessage(long orderId, long businessId, String productName,
int quantity, double unitPrice) {
this.orderId = orderId;
this.businessId = businessId;
this.productName = productName;
this.quantity = quantity;
this.unitPrice = unitPrice;
}
public long getOrderId() { return orderId; }
public long getBusinessId() { return businessId; }
public String getProductName() { return productName; }
public int getQuantity() { return quantity; }
public double getUnitPrice() { return unitPrice; }
@Override
public String toString() {
return "OrderMessage{orderId=" + orderId + ", businessId=" + businessId
+ ", productName='" + productName + "', quantity=" + quantity
+ ", unitPrice=" + unitPrice + "}";
}
}
import java.io.Serializable;
/** 输出消息:加工完成后发送到 flink_out 主题 */
public class OutMessage implements Serializable {
private final long orderId;
private final long channelId;
private final String productName;
private final int quantity;
private final double totalPrice;
public OutMessage(long orderId, long channelId, String productName,
int quantity, double totalPrice) {
this.orderId = orderId;
this.channelId = channelId;
this.productName = productName;
this.quantity = quantity;
this.totalPrice = totalPrice;
}
public long getOrderId() { return orderId; }
public long getChannelId() { return channelId; }
public String getProductName() { return productName; }
public int getQuantity() { return quantity; }
public double getTotalPrice() { return totalPrice; }
@Override
public String toString() {
return "OutMessage{orderId=" + orderId + ", channelId=" + channelId
+ ", productName='" + productName + "', quantity=" + quantity
+ ", totalPrice=" + totalPrice + "}";
}
}
6. LookupEventBus 事件总线
LookupEventBus 是 CSRCDSEPC 的核心解耦机制,按 key 发布-订阅。这里定义两条总线:订单处理总线和输出总线:
import akka.actor.ActorRef;
import akka.event.japi.LookupEventBus;
/** 订单处理总线:按业务线 ID 投递 */
public class OrderProcessBus extends LookupEventBus<OrderMessage, ActorRef, Long> {
@Override
public int mapSize() { return 128; }
@Override
public int compareSubscribers(ActorRef a, ActorRef b) {
return a.compareTo(b);
}
@Override
public Long classify(OrderMessage event) {
return event.getBusinessId();
}
@Override
public void publish(OrderMessage event, ActorRef subscriber) {
subscriber.tell(event, ActorRef.noSender());
}
}
import akka.actor.ActorRef;
import akka.event.japi.LookupEventBus;
/** 输出总线:按输出通道 ID 投递 */
public class OutputBus extends LookupEventBus<OutMessage, ActorRef, Long> {
private static final OutputBus INSTANCE = new OutputBus();
private OutputBus() {}
public static OutputBus getInstance() { return INSTANCE; }
@Override
public int mapSize() { return 128; }
@Override
public int compareSubscribers(ActorRef a, ActorRef b) {
return a.compareTo(b);
}
@Override
public Long classify(OutMessage event) {
return event.getChannelId();
}
@Override
public void publish(OutMessage event, ActorRef subscriber) {
subscriber.tell(event, ActorRef.noSender());
}
}
7. Kafka 消费与模拟发送
KafkaConsumerActor 模拟从 Kafka 消费消息,OutputPipeline 模拟 KafkaProducer 发送到 flink_out 主题:
import akka.actor.AbstractActor;
import akka.actor.ActorRef;
import akka.actor.Props;
/** 模拟 Kafka 消费者:从 Kafka 拉取消息并发布到订单处理总线 */
public class KafkaConsumerActor extends AbstractActor {
private final ActorRef processBus;
private final long businessId;
public KafkaConsumerActor(ActorRef processBus, long businessId) {
this.processBus = processBus;
this.businessId = businessId;
}
public static Props props(ActorRef processBus, long businessId) {
return Props.create(KafkaConsumerActor.class, processBus, businessId);
}
@Override
public Receive createReceive() {
return receiveBuilder()
.match(OrderMessage.class, msg -> {
// 模拟从 Kafka 消费到消息,发布到订单处理总线
processBus.tell(msg, getSelf());
})
.build();
}
}
import akka.actor.AbstractActor;
import akka.actor.Props;
/** 输出管道:模拟 KafkaProducer 发送到 flink_out 主题 */
public class OutputPipeline extends AbstractActor {
public static Props props() {
return Props.create(OutputPipeline.class);
}
@Override
public Receive createReceive() {
return receiveBuilder()
.match(OutMessage.class, msg -> {
// 模拟 KafkaProducer.send() 发送到 flink_out 主题
System.out.println("[KafkaProducer] 发送到 flink_out: " + msg);
})
.build();
}
}
8. Stepper 处理链
每个 Planner 内部串联一条 Stepper 链,依次完成校验、取价、组装和输出:
import akka.actor.AbstractActor;
import akka.actor.ActorRef;
import akka.actor.Props;
/** 校验 Stepper:校验订单并拆分商品 */
public class ValidateStepper extends AbstractActor {
private final ActorRef next;
public ValidateStepper(ActorRef next) {
this.next = next;
}
public static Props props(ActorRef next) {
return Props.create(ValidateStepper.class, next);
}
@Override
public Receive createReceive() {
return receiveBuilder()
.match(OrderMessage.class, msg -> {
if (msg.getQuantity() <= 0) {
System.out.println("[Validate] 订单 " + msg.getOrderId() + " 数量非法,丢弃");
return;
}
System.out.println("[Validate] 订单 " + msg.getOrderId() + " 校验通过");
next.forward(msg, getContext());
})
.build();
}
}
import akka.actor.AbstractActor;
import akka.actor.ActorRef;
import akka.actor.Props;
/** 取价 Stepper:模拟远程取价 */
public class PriceStepper extends AbstractActor {
private final ActorRef next;
public PriceStepper(ActorRef next) {
this.next = next;
}
public static Props props(ActorRef next) {
return Props.create(PriceStepper.class, next);
}
@Override
public Receive createReceive() {
return receiveBuilder()
.match(OrderMessage.class, msg -> {
// 模拟远程取价:单价上浮 10% 作为加工结果
double processedPrice = msg.getUnitPrice() * 1.10;
OrderMessage processed = new OrderMessage(
msg.getOrderId(), msg.getBusinessId(),
msg.getProductName(), msg.getQuantity(), processedPrice);
System.out.println("[Price] 订单 " + msg.getOrderId()
+ " 取价完成: " + processedPrice);
next.forward(processed, getContext());
})
.build();
}
}
import akka.actor.AbstractActor;
import akka.actor.ActorRef;
import akka.actor.Props;
/** 结果 Stepper:组装输出消息 */
public class ResultStepper extends AbstractActor {
private final ActorRef next;
private final long outputChannelId;
public ResultStepper(ActorRef next, long outputChannelId) {
this.next = next;
this.outputChannelId = outputChannelId;
}
public static Props props(ActorRef next, long outputChannelId) {
return Props.create(ResultStepper.class, next, outputChannelId);
}
@Override
public Receive createReceive() {
return receiveBuilder()
.match(OrderMessage.class, msg -> {
double totalPrice = msg.getUnitPrice() * msg.getQuantity();
OutMessage out = new OutMessage(
msg.getOrderId(), outputChannelId,
msg.getProductName(), msg.getQuantity(), totalPrice);
System.out.println("[Result] 订单 " + msg.getOrderId()
+ " 组装完成: " + out);
next.forward(out, getContext());
})
.build();
}
}
package com.epc.web.akka;
import akka.actor.*;
import akka.japi.pf.ReceiveBuilder;
/**
* 输出节点(≈ CSRCDSEPC 的 OutputStepper)。
*
* <p>流水线末端。收到 ResultStepper 组装好的 {@link OutMessage} 后,
* 直接 {@code OutputBus.publish(...)} 发到输出总线。</p>
*
* <p>它继承 OrderStepper 仅为复用基类的日志辅助;但实际处理的是 OutMessage
* (不是 OrderMessage),所以覆写 createReceive 直接匹配 OutMessage。</p>
*/
public class OutputStepper extends OrderStepper {
private OutputStepper() {
super(null); // 末端,没有下一级
}
public static Props props() {
return Props.create(OutputStepper.class);
}
@Override
public Receive createReceive() {
return ReceiveBuilder.create()
.match(OutMessage.class, msg -> {
getContext().getSystem().log().info(
"[" + tag() + "] 发送到输出总线: " + msg);
// 发到全局输出总线(≈ OutputServiceV2.outputBus.publish)
OutputBus.getInstance().publish(msg);
})
.build();
}
/** OrderStepper 模板要求的抽象方法,本节点不用 */
@Override
protected void doProcess(OrderMessage msg) {
// no-op:OutMessage 在 createReceive 里直接处理
}
}
9. Planner 路由器
OrderPlanner 使用 RoundRobinPool 路由器,从配置读取动态扩缩参数,每个 routee 内部组装一条 Stepper 链:
import akka.actor.AbstractActor;
import akka.actor.ActorRef;
import akka.actor.Props;
import akka.routing.FromConfig;
/** 订单规划器:RoundRobinPool 路由器,动态扩缩容 */
public class OrderPlanner extends AbstractActor {
private final ActorRef outputBus;
private final long outputChannelId;
public OrderPlanner(ActorRef outputBus, long outputChannelId) {
this.outputBus = outputBus;
this.outputChannelId = outputChannelId;
}
public static Props props(ActorRef outputBus, long outputChannelId) {
return Props.create(OrderPlanner.class, outputBus, outputChannelId);
}
@Override
public Receive createReceive() {
return receiveBuilder()
.match(OrderMessage.class, msg -> {
// 每个 routee 内部组装 Stepper 链
ActorRef outputStepper = getContext().actorOf(
OutputStepper.props(outputBus), "outputStepper");
ActorRef resultStepper = getContext().actorOf(
ResultStepper.props(outputStepper, outputChannelId), "resultStepper");
ActorRef priceStepper = getContext().actorOf(
PriceStepper.props(resultStepper), "priceStepper");
ActorRef validateStepper = getContext().actorOf(
ValidateStepper.props(priceStepper), "validateStepper");
// 将消息转发给链首
validateStepper.forward(msg, getContext());
})
.build();
}
/** 创建带路由器的 Planner */
public static Props routerProps(ActorRef outputBus, long outputChannelId) {
return FromConfig.getInstance().props(
Props.create(OrderPlanner.class, outputBus, outputChannelId));
}
}
10. 入口类 DemoApp
DemoApp 装配两个 ActorSystem、两条 LookupEventBus,灌入若干订单,观察消息如何经过 Planner → Stepper 链 → 输出总线 → 模拟 Kafka:
/*
* Akka Actor 可运行示例(入口类)
* ============================================================================
* 本示例复刻 CSRCDSEPC dataset 模块(DatasetServiceV2 / Planner / OutputServiceV2)
* 用到的全部 Akka 核心能力,但剥离了 Spring / Kafka / DB / Groovy,
* 用一个"订单处理"的小场景自洽地跑通,方便阅读与运行。
*
* 涵盖的 Akka 知识点:
* 1. ActorSystem —— Actor 的容器/运行时(这里建了两个独立 system)
* 2. AbstractActor —— 用 Java 写 Actor(createReceive 定义消息处理)
* 3. Props —— Actor 的"工厂蓝图"
* 4. actorOf —— 创建 Actor 实例
* 5. Router + RoundRobinPool —— 把 1 个逻辑 Actor 路由到 N 个 routee 并行处理
* 6. FromConfig —— 路由器参数从 application.conf 读取(动态扩缩容)
* 7. LookupEventBus —— 按 key 发布-订阅的事件总线(CSRCDSEPC 的核心解耦机制)
* 8. PoisonPill —— 优雅关闭 Actor
* 9. forward / tell —— Actor 之间异步通信
*
* 数据流(和 CSRCDSEPC 的 dataset 流程一一对应):
*
* 订单消息 OrderMessage
* │ DemoApp 逐条 publish(≈ DatasetServiceV2.process)
* ▼
* OrderProcessBus (LookupEventBus, key=Long 业务线ID)
* │ 只投递给"订阅了这条业务线"的 router(≈ sourceBus 按 sourceId 投递)
* ▼
* OrderPlanner (Router: RoundRobinPool 2~5 个 routee 动态扩缩)
* │ 每个 Planner 内部再串一条 actor 链(≈ Planner.preStart 组装 stepper 链)
* ├─▶ ValidateStepper ──▶ PriceStepper ──▶ ResultStepper ──▶ OutputStepper
* │ (校验+拆分) (模拟远程取价) (组装结果) (发到输出总线)
* └────────────────────────────────────────────────────────────┘
* │
* ▼
* OutputBus (LookupEventBus, key=Long 输出通道ID)
* │ 订阅该通道的 router
* ▼
* OutputPipeline (Router) ──▶ 模拟 KafkaProducer.send()*/
package com.epc.web.akka;
import akka.actor.ActorRef;
import akka.actor.ActorSystem;
import akka.actor.PoisonPill;
import akka.event.japi.LookupEventBus;
import akka.routing.Broadcast;
import com.typesafe.config.Config;
import com.typesafe.config.ConfigFactory;
import java.util.concurrent.TimeUnit;
/**
* 演示入口:装配两个 ActorSystem + 两条 LookupEventBus,灌入若干订单,
* 观察消息如何经过 Planner → Stepper 链 → 输出总线 → 模拟 Kafka。
* > 技术栈与主项目一致:**JDK 1.8 + Akka 2.5.11 + Scala 2.11**。
* > 文件位于 `src/test`,走 Maven 测试类路径即可运行(无需真实 Kafka / DB)。
*/
public class DemoApp {
/**
* 业务线ID:订单处理总线的路由键(≈ CSRCDSEPC 的方案ID / sourceId)。
* 所有演示订单都归属这一条业务线,所以会被同一个 Planner router 接收。
*/
public static final long ORDER_BUSINESS_ID = 1001L;
/**
* 输出通道ID:输出总线的路由键(≈ CSRCDSEPC 的 pipelineId)。
* public 以便 ResultStepper 组装 OutMessage 时引用。
*/
public static final long OUTPUT_CHANNEL_ID = 2001L;
/**
* 输出总线(≈ CSRCDSEPC 里 OutputServiceV2.outputBus 是 static)。
* 用单例:OutputStepper 在另一个 ActorSystem 里也要往它 publish,
* 必须保证 DemoApp 订阅和 Stepper 发布用的是同一个 bus 实例。
*/
public static final OutputBus OUTPUT_BUS = OutputBus.getInstance();
public static void main(String[] args) throws Exception {
// ── 1. 加载配置 ──────────────────────────────────────────────
// ConfigFactory.load() 会读 classpath 上的 application.conf。
// 这里用 demo.conf 里的两段(OrderDemoProcessSystem / OrderDemoOutputSystem)
// 给两个 system 各配一个带"动态扩缩 resizer"的 RoundRobinPool 路由器。
Config config = ConfigFactory.load("demo").withFallback(ConfigFactory.load());
System.out.println("===== Akka Actor 演示开始 =====");
// 重置 MockKafka 计数:本次演示预期 5 个订单 × 各拆 2 件 = 10 条输出
int expectedOutputs = 5 * 2;
MockKafka.reset(expectedOutputs);
// ── 2. 第一条 ActorSystem:订单处理(≈ DatasetProcessSystem)──
// 名字必须和 demo.conf 里的段名一致,FromConfig 才能找到 deployment 路由配置。
ActorSystem orderSystem = ActorSystem.create("OrderDemoProcessSystem", config);
try {
// 订单处理总线:消息按"业务线ID"分类,投递给订阅了该业务线的 router。
LookupEventBus<OrderMessage, Subscriber, Long> orderBus = new OrderProcessBus();
// 创建 Planner router 并"订阅"业务线 1001 的订单。
// (≈ EventBusSubscription.subscribe:建 router → Subscriber.of → bus.subscribe)
ActorRef plannerRouter = orderSystem.actorOf(
OrderPlanner.props(),
"orderPlanner");
orderBus.subscribe(Subscriber.of(plannerRouter, "orderPlanner"), ORDER_BUSINESS_ID);
System.out.println("[init] orderPlanner router 已创建并订阅业务线 " + ORDER_BUSINESS_ID);
// 输出 ActorSystem(≈ OutputSystem)
ActorSystem outputSystem = ActorSystem.create("OrderDemoOutputSystem", config);
try {
// 输出管道 router,订阅输出通道 2001
ActorRef pipelineRouter = outputSystem.actorOf(
OutputPipeline.props(),
"outputPipeline");
OUTPUT_BUS.subscribe(Subscriber.of(pipelineRouter, "outputPipeline"), OUTPUT_CHANNEL_ID);
System.out.println("[init] outputPipeline router 已创建并订阅输出通道 " + OUTPUT_CHANNEL_ID);
// ── 3. 灌入 5 条订单,观察整条流水线 ──────────────────
// 每条 publish 都会进入 orderBus.classify() 取 businessId,
// 再 tell 给上面订阅了 1001 的 orderPlanner router,
// router 再把消息轮询派发给某个 Planner routee。
for (int i = 1; i <= 5; i++) {
orderBus.publish(new OrderMessage(ORDER_BUSINESS_ID, "ord-" + i, "商品" + i, i * 10L, 2));
System.out.println("[main] 已 publish 订单 ord-" + i + "(businessId=" + ORDER_BUSINESS_ID + ")");
}
// ── 4. 等待全部输出完成(CountDownLatch,见 MockKafka)──
boolean done = MockKafka.waitForDrain(expectedOutputs, 20);
System.out.println("[main] 等待结果: " + (done ? "全部输出完成" : "超时(仍有消息在途)"));
} finally {
// 优雅关闭输出 system:广播 PoisonPill(≈ OutputServiceV2.terminateActors)
outputSystem.actorSelection("/user/*")
.tell(new Broadcast(PoisonPill.getInstance()), ActorRef.noSender());
outputSystem.terminate();
}
} finally {
// 优雅关闭订单处理 system
orderSystem.actorSelection("/user/*")
.tell(new Broadcast(PoisonPill.getInstance()), ActorRef.noSender());
orderSystem.terminate();
}
// 最后打印 MockKafka 收到的全部消息
MockKafka.printAll();
System.out.println("===== Akka Actor 演示结束 =====");
// 给 JVM 一点时间收尾(测试环境下 main 线程结束即可退出)
Thread.sleep(TimeUnit.SECONDS.toMillis(1));
}
}
转载自 CSDN-专业IT技术社区
原文链接:https://blog.csdn.net/clz1314521/article/details/165122465



