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/ 自动加载
✔ 无需重写主程序
✔ 支持扩展、隔离、版本化

更多推荐