仿照Flink基于SPI机制实现插件化消息发送服务
1️⃣ 需求与整体架构
为了做到像 Flink 一样的可插拔:
-
核心服务(不可变)
提供统一接口 + 插件加载框架 -
插件(可变)
每种实现(短信/电话/邮件)单独开发成一个 jar 包 -
运行时扫描 + ServiceLoader 加载
主服务无需编译时依赖实现类
架构图(文本图示):

2️⃣ 核心 SPI 接口(core-api.jar)
所有插件都必须实现这个接口。
package com.example.messaging.spi;
import java.util.Map;
public interface MessageSender {
String name(); // 插件唯一标识
String type(); // 消息类型:sms/phone/email/webhook
void init(Map<String, String> config) throws Exception;
SendResult send(Message message) throws Exception;
void close() throws Exception;
}
Message & SendResult
public class Message {
private final String to;
private final String body;
private final Map<String,String> attrs;
// constructor / getter
}
public class SendResult {
private final boolean success;
private final String code;
private final String errorMessage;
// constructor / getter
}
3️⃣ 插件约定(META-INF/services)
所有消息实现必须在 jar 包中加入:
META-INF/services/com.example.messaging.spi.MessageSender
内容示例:
com.vendor.sms.SmsSender
Java 会使用 ServiceLoader 自动发现和实例化。
4️⃣ 插件加载器 PluginManager(Flink 核心方式)
核心能力:
-
遍历
plugins/目录 -
每个 jar 用独立 ClassLoader 加载
-
ServiceLoader 自动扫描实现类
-
注册到内存表 name/type → 实例
注意:示例使用 parent-first,生产可换成 Flink 的 parent-last classloader。
PluginManager.java
package com.example.messaging.core;
import com.example.messaging.spi.MessageSender;
import java.io.File;
import java.net.URL;
import java.net.URLClassLoader;
import java.util.*;
public class PluginManager {
private final File pluginsDir;
private final Map<String, PluginHolder> loaded = new HashMap<>();
public PluginManager(File pluginsDir) {
this.pluginsDir = pluginsDir;
}
public void loadAll() throws Exception {
if (!pluginsDir.exists()) return;
File[] jars = pluginsDir.listFiles(
f -> f.isFile() && f.getName().endsWith(".jar")
);
if (jars == null) return;
for (File jar : jars) {
loadSingleJar(jar);
}
}
private void loadSingleJar(File jar) throws Exception {
URLClassLoader cl = new URLClassLoader(
new URL[]{jar.toURI().toURL()},
this.getClass().getClassLoader()
);
ServiceLoader<MessageSender> sl = ServiceLoader.load(
MessageSender.class, cl
);
for (MessageSender sender : sl) {
loaded.put(sender.name(),
new PluginHolder(sender.name(), sender.type(), sender, cl)
);
System.out.println("[Plugin] Loaded: " + sender.name());
}
}
public Optional<MessageSender> getByName(String name) {
PluginHolder h = loaded.get(name);
return h == null ? Optional.empty() : Optional.of(h.sender);
}
public List<MessageSender> getByType(String type) {
List<MessageSender> out = new ArrayList<>();
for (PluginHolder h : loaded.values()) {
if (h.type.equals(type)) out.add(h.sender);
}
return out;
}
private static class PluginHolder {
String name, type;
MessageSender sender;
ClassLoader cl;
PluginHolder(String n, String t, MessageSender s, ClassLoader c) {
name = n; type = t; sender = s; cl = c;
}
}
}
5️⃣ 消息服务入口:MessageService
业务调用只需要调用它。
public class MessageService {
private final PluginManager pm;
public MessageService(File pluginDir) throws Exception {
this.pm = new PluginManager(pluginDir);
pm.loadAll();
}
public SendResult sendByName(String name, Message msg, Map<String,String> cfg) throws Exception {
MessageSender s = pm.getByName(name)
.orElseThrow(() -> new RuntimeException("Sender not found: " + name));
s.init(cfg);
return s.send(msg);
}
public SendResult sendByType(String type, Message msg, Map<String,String> cfg) throws Exception {
List<MessageSender> list = pm.getByType(type);
if (list.isEmpty())
throw new RuntimeException("No sender for type " + type);
MessageSender s = list.get(0); // 简单策略
s.init(cfg);
return s.send(msg);
}
}
6️⃣ 插件实现示例(短信 SmsSender)
SmsSender.java
package com.vendor.sms;
import com.example.messaging.spi.*;
import java.util.Map;
public class SmsSender implements MessageSender {
@Override
public String name() {
return "sms-vendor-v1";
}
@Override
public String type() {
return "sms";
}
@Override
public void init(Map<String,String> config) {
// 初始化 API Key / HTTP 客户端
}
@Override
public SendResult send(Message message) throws Exception {
// 调用第三方短信 API
System.out.println("Send SMS to " + message.getTo());
return new SendResult(true, "200", null);
}
@Override
public void close() {}
}
META-INF/services/com.example.messaging.spi.MessageSender
com.vendor.sms.SmsSender
7️⃣ 完整目录结构(可直接照着搭项目)
messaging-service/
├─ app/
│ └─ app.jar (含 PluginManager/MessageService)
├─ core-api/
│ └─ core-api.jar
├─ plugins/
│ ├─ sms-sender-1.0.jar
│ └─ phone-sender-1.0.jar
└─ conf/
└─ app.conf
8️⃣ 启动到发送的流程图(文本版)
[Start]
│
▼
读取 conf 配置
│
▼
PluginManager 扫描 plugins/ 目录
│
▼
为每个 jar 创建 ClassLoader
│
▼
ServiceLoader 加载 MessageSender 实现
│
▼
注册 sender(name/type → 实例)
│
▼
等待业务调用 send()
│
──────────────────────────────────────────
│ 调用 sendByType("sms") or sendByName() │
──────────────────────────────────────────
▼
匹配对应 sender
│
▼
执行 sender.init()
│
▼
执行 sender.send()
│
▼
返回 SendResult
▼
[End]
9️⃣ 如何新增一个消息类型?只需 3 步
Step 1:实现 MessageSender
public class EmailSender implements MessageSender {...}
Step 2:添加 META-INF/services
META-INF/services/com.example.messaging.spi.MessageSender
内容:
com.xxx.EmailSender
Step 3:打成 jar → 丢到 plugins/ → 重启服务
无需修改主程序。完全像 Flink 插件一样。
🔟 总结
本文实现了一个高度模块化、可升级、可热插拔的消息服务框架,其核心思想来自:
-
Java SPI(ServiceLoader)
-
Flink Plugin 机制(ClassLoader 隔离)
最终实现:
✔ 每种消息独立成 jar
✔ 放入 plugins/ 自动加载
✔ 无需重写主程序
✔ 支持扩展、隔离、版本化
更多推荐
所有评论(0)