# FileBasedMessageQueue **Repository Path**: love47/file-based-message-queue ## Basic Information - **Project Name**: FileBasedMessageQueue - **Description**: 基于文件流的伪消息队列 - **Primary Language**: Unknown - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 2 - **Forks**: 0 - **Created**: 2026-09-08 - **Last Updated**: 2026-09-20 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # FileMQ - 基于 FTP 单向同步的类消息队列 FileMQ 是一个基于文件系统实现的消息队列,适用于网络隔离环境,通过 FTP 单向同步实现消息传递,未来可平滑迁移到标准 MQ(Kafka/RabbitMQ)。 ## 特性 - ✅ **高可靠性**:fsync 保证消息持久化,原子 rename 避免半包消息 - ✅ **幂等性保证**:消息 ID 去重,数据库记录消费日志 - ✅ **顺序保证**:单分区严格有序,多分区按 Key 有序 - ✅ **自动重试**:延迟重试 + 死信队列 - ✅ **Spring Boot 集成**:提供 Starter,开箱即用 - ✅ **注解驱动**:@FileMQListener 简化消费者开发 - ✅ **平滑迁移**:抽象接口层,未来可切换到标准 MQ ## 技术栈 - Java 8+ - Spring Boot 2.7+ - MyBatis-Plus 3.5+ - MySQL 8.0+ ## 快速开始 ### 1. 引入依赖 ```xml com.filemq filemq-spring-boot-starter 1.0.0 ``` ### 2. 配置数据库 ```yaml spring: datasource: url: jdbc:mysql://localhost:3306/filemq?useUnicode=true&characterEncoding=utf8&useSSL=false&serverTimezone=Asia/Shanghai username: root password: root ``` ### 3. 启用 FileMQ ```java @SpringBootApplication @EnableFileMQ public class Application { public static void main(String[] args) { SpringApplication.run(Application.class, args); } } ``` ### 4. 发送消息 ```java @Autowired private FileSystemProducer producer; public void sendMessage() { MessageFuture future = producer.send("order", "{\"orderId\":\"123\"}"); if (future.isSuccess()) { System.out.println("Sent: " + future.getMsgId()); } } ``` ### 5. 接收消息 ```java @Component public class OrderListener { @FileMQListener(topic = "order") public void handleOrder(Message message) { System.out.println("Received: " + message.getBody()); } } ``` ## 项目结构 ``` filemq-parent/ ├── filemq-spring-boot-starter/ # 核心启动器 │ └── src/main/java/com/filemq/ │ ├── annotation/ # @FileMQListener │ ├── autoconfigure/ # Spring Boot 自动配置 │ ├── core/ # 核心抽象接口 │ ├── producer/ # 生产者实现 │ ├── consumer/ # 消费者实现 │ ├── entity/ # JPA/MyBatis 实体 │ └── mapper/ # MyBatis Mapper ├── filemq-samples/ # 示例代码 └── USAGE.md # 详细使用手册 ``` ## 核心概念 ### 生产端(Outbox) ``` outbox/ ├── .staging/ # 临时写入目录 └── ready/ # 就绪目录,FTP 从这里读取 ``` 发送流程:staging → fsync → 原子 rename → ready ### 消费端(Inbox) ``` inbox/ ├── ready/ # FTP 写入这里 ├── locked/ # 处理中(原子 rename 实现锁) ├── completed/ # 处理完成(按日期归档) ├── failed/ # 死信队列 └── delayed/ # 延迟重试 ``` 消费流程:ready → 原子 rename → locked → 处理 → completed/failed/delayed ### 消息协议 文件名格式: ``` {timestamp}_{seq}_{msgId}_{topic}_{partition}_{checksum}.msg ``` 消息文件格式: ``` [Header] version: 1.0 msgid: xxx topic: order ... [Body] {"orderId":"123"} ``` ## 详细文档 请查看 [USAGE.md](./USAGE.md) 获取完整的使用手册: - 配置说明 - 发送/接收消息的多种方式 - 数据库初始化 - 最佳实践 - 常见问题 ## 构建和运行 ```bash # 编译整个项目 mvn clean install # 运行示例 cd filemq-samples mvn spring-boot:run ``` ## FTP 同步配置 第三方 FTP 同步工具需要配置: - 源目录:`outbox/ready` - 目标目录:`inbox/ready` - 策略:单向同步,传输完成后删除源文件 示例 rsync 脚本: ```bash rsync -av --remove-source-files /path/to/outbox/ready/ user@remote:/path/to/inbox/ready/ ``` ## 未来迁移到标准 MQ 当网络条件允许时,可以平滑迁移: 1. 保持业务代码不变(使用 MessageProducer/MessageConsumer 接口) 2. 替换实现为 Kafka/RabbitMQ 版本 3. 消息格式保持兼容 ## 许可证 MIT License ## 联系方式 如有问题,请提交 Issue。