Java 程序员第 42 阶段17:文档智能解析审核,大模型实现合同摘要与合规校验,分布式文档处理架构设计
目录
1. 章节简介
2. 分布式文档处理架构概述
3. 消息队列集成(Kafka/RabbitMQ)
4. 文档处理任务调度
5. 失败重试与幂等处理
6. 实际代码示例:消息生产者/消费者实现
7. 实际代码示例:任务调度服务
8. 实际代码示例:分布式锁实现
9. 章节总结
1. 章节简介
本章节将深入探讨文档智能解析审核系统中的分布式架构设计。在实际生产环境中,文档处理系统面临高并发、大数据量的挑战,单机部署模式已无法满足业务需求。本章节将详细介绍如何通过分布式架构、消息队列集成、任务调度系统来实现高性能、高可用的文档处理平台。
1.1 章节内容概述
本章节涵盖以下核心技术点:
- **分布式架构设计**:阐述整体架构设计理念,包括服务拆分、节点规划、数据分片等
- **消息队列集成**:深入讲解Kafka和RabbitMQ在系统中的集成方案和最佳实践
- **任务调度系统**:介绍分布式任务调度的实现原理和主流框架
- **失败重试机制**:讲解如何保证任务在分布式环境下的可靠执行
- **幂等性保证**:阐述如何设计幂等操作,防止重复处理
1.2 学习目标
通过本章节的学习,您将掌握:
1. 理解分布式文档处理系统的架构设计原理
2. 掌握Kafka/RabbitMQ消息队列的配置和集成方法
3. 熟悉分布式任务调度框架的使用
4. 学会设计高可用的失败重试和幂等处理机制
5. 能够编写生产级的消息生产者/消费者代码
2. 分布式文档处理架构概述

2.1 架构设计理念
分布式文档处理架构遵循以下核心设计理念:
#### 2.1.1 服务拆分原则
系统采用微服务架构思想,将整体业务拆分为多个独立的服务:
- **API网关服务**:统一入口,负责请求路由、负载均衡、认证授权
- **文档解析服务**:核心业务服务,负责文档内容的解析和处理
- **任务调度服务**:管理系统中的定时任务和异步任务
- **通知服务**:负责向用户发送处理结果通知
#### 2.1.2 可扩展性设计
系统支持水平和垂直两个方向的扩展:
|
扩展方向 |
实现方式 |
适用场景 |
|
水平扩展 |
增加处理节点 |
高并发、大数据量 |
|
垂直扩展 |
升级硬件配置 |
单机性能瓶颈 |
#### 2.1.3 高可用设计
系统通过以下机制保证高可用:
- 多副本部署:关键服务至少部署2个副本
- 故障自动转移:节点故障时自动切换到健康节点
- 熔断降级:服务不可用时返回降级响应
2.2 系统架构图
`
┌─────────────────────────────────────────────────────────────┐
│ API Gateway │
│ (负载均衡 / 路由) │
└────────────────────────┬──────────────────────────────────┘
│
┌───────────────┼───────────────┐
│ │ │
▼ ▼ ▼
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ Kafka │ │ RabbitMQ │ │ 存储层 │
│ 消息队列集群 │ │ 消息队列集群 │ │ (MySQL/ │
│ │ │ │ │ MongoDB/ │
│ Topic: │ │ Exchange: │ │ Redis) │
│ doc-parse │ │ doc.exchange│ │ │
│ doc-review │ │ │ └─────────────┘
│ doc-notify │ │ Queue: │
│ doc-dlq │ │ doc.parse │
└─────────────┘ │ doc.review │
│ doc.dlq │
└─────────────┘
│
┌───────────────┼───────────────┐
│ │ │
▼ ▼ ▼
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ Processing │ │ Processing │ │ Processing │
│ Node 1 │ │ Node 2 │ │ Node N │
│ Worker │ │ Worker │ │ Worker │
└─────────────┘ └─────────────┘ └─────────────┘
`
2.3 核心组件介绍
#### 2.3.1 API网关
API网关是系统的统一入口,主要功能包括:
- **负载均衡**:将请求均匀分发到后端服务节点
- **路由转发**:根据请求路径转发到对应的微服务
- **认证授权**:验证用户身份和权限
- **限流熔断**:防止系统过载,保证稳定性
#### 2.3.2 消息队列集群
消息队列承担着异步通信和系统解耦的重任:
- **Kafka**:适用于高吞吐量的日志收集和消息推送场景
- **RabbitMQ**:适用于需要灵活路由和复杂消费模式的场景
#### 2.3.3 处理节点集群
处理节点是实际执行业务逻辑的工作节点:
- 采用无状态设计,可随时水平扩展
- 通过消息队列获取任务,处理完成后更新任务状态
- 支持资源隔离,不同类型的任务可以分配到不同的节点池
#### 2.3.4 存储层
存储层采用多元存储策略:
|
存储类型 |
使用场景 |
代表产品 |
|
关系型存储 |
结构化业务数据 |
MySQL |
|
文档型存储 |
合同原文、元数据 |
MongoDB |
|
缓存存储 |
会话、锁、热点数据 |
Redis |
|
对象存储 |
原始文档文件 |
MinIO |
3. 消息队列集成(Kafka/RabbitMQ)

3.1 消息队列选型考量
在分布式系统中,消息队列扮演着至关重要的角色。选型时需要考虑以下因素:
|
考量因素 |
Kafka |
RabbitMQ |
|
吞吐量 |
百万级/秒 |
万级/秒 |
|
延迟 |
毫秒级 |
微秒级 |
|
消息持久化 |
支持 |
支持 |
|
消息回溯 |
支持 |
不支持 |
|
路由模式 |
Topic/Partition |
Exchange/Queue |
|
集群部署 |
原生支持 |
需要 federation 插件 |
本系统根据不同业务场景,采用双队列组合方案:
- **Kafka**:用于高吞吐量的文档解析任务分发
- **RabbitMQ**:用于需要精确路由的通知服务
3.2 Kafka配置与集成
#### 3.2.1 Topic设计
```properties
Kafka Topic 配置
doc-parse partitions=6 replication-factor=3
doc-review partitions=6 replication-factor=3
doc-notify partitions=3 replication-factor=2
doc-dlq partitions=3 replication-factor=2
`
#### 3.2.2 生产者配置
```java
@Configuration
public class KafkaProducerConfig {
@Value("${spring.kafka.bootstrap-servers}")
private String bootstrapServers;
@Bean
public ProducerFactory<String, DocumentMessage> producerFactory() {
Map<String, Object> configProps = new HashMap<>();
configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
// 可靠性配置
configProps.put(ProducerConfig.ACKS_CONFIG, "all");
configProps.put(ProducerConfig.RETRIES_CONFIG, 3);
configProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
// 性能优化
configProps.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384);
configProps.put(ProducerConfig.LINGER_MS_CONFIG, 5);
configProps.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432);
return new DefaultKafkaProducerFactory<>(configProps);
}
@Bean
public KafkaTemplate<String, DocumentMessage> kafkaTemplate() {
return new KafkaTemplate<>(producerFactory());
}
}
`
#### 3.2.3 消费者配置
```java
@Configuration
@EnableKafka
public class KafkaConsumerConfig {
@Value("${spring.kafka.bootstrap-servers}")
private String bootstrapServers;
@Value("${spring.kafka.consumer.group-id}")
private String groupId;
@Bean
public ConsumerFactory<String, DocumentMessage> consumerFactory() {
Map<String, Object> props = new HashMap<>();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
// 自动提交偏移量
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
// 拉取策略
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 100);
props.put(ConsumerConfig.FETCH_MIN_BYTES_CONFIG, 1024);
props.put(ConsumerConfig.FETCH_MAX_WAIT_MS_CONFIG, 500);
return new DefaultKafkaConsumerFactory<>(props,
new StringDeserializer(),
new JsonDeserializer<>(DocumentMessage.class));
}
@Bean
public ConcurrentKafkaListenerContainerFactory<String, DocumentMessage>
kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, DocumentMessage> factory =
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());
factory.setConcurrency(3);
factory.getContainerProperties().setAckMode(AckMode.MANUAL_IMMEDIATE);
return factory;
}
}
`
3.3 RabbitMQ配置与集成
#### 3.3.1 Exchange和Queue配置
```java
@Configuration
public class RabbitMQConfig {
public static final String DOC_EXCHANGE = "doc.exchange";
public static final String DOC_PARSE_QUEUE = "doc.parse.queue";
public static final String DOC_REVIEW_QUEUE = "doc.review.queue";
public static final String DOC_NOTIFY_QUEUE = "doc.notify.queue";
public static final String DOC_DLQ = "doc.dlq";
public static final String ROUTING_KEY_PARSE = "doc.parse";
public static final String ROUTING_KEY_REVIEW = "doc.review";
public static final String ROUTING_KEY_NOTIFY = "doc.notify";
@Bean
public DirectExchange docExchange() {
return new DirectExchange(DOC_EXCHANGE, true, false);
}
@Bean
public Queue docParseQueue() {
return QueueBuilder.durable(DOC_PARSE_QUEUE)
.withArgument("x-dead-letter-exchange", "")
.withArgument("x-dead-letter-routing-key", DOC_DLQ)
.build();
}
@Bean
public Queue docReviewQueue() {
return QueueBuilder.durable(DOC_REVIEW_QUEUE)
.withArgument("x-dead-letter-exchange", "")
.withArgument("x-dead-letter-routing-key", DOC_DLQ)
.build();
}
@Bean
public Queue docNotifyQueue() {
return QueueBuilder.durable(DOC_NOTIFY_QUEUE)
.withArgument("x-dead-letter-exchange", "")
.withArgument("x-dead-letter-routing-key", DOC_DLQ)
.build();
}
@Bean
public Queue docDeadLetterQueue() {
return QueueBuilder.durable(DOC_DLQ).build();
}
@Bean
public Binding docParseBinding() {
return BindingBuilder.bind(docParseQueue())
.to(docExchange())
.with(ROUTING_KEY_PARSE);
}
@Bean
public Binding docReviewBinding() {
return BindingBuilder.bind(docReviewQueue())
.to(docExchange())
.with(ROUTING_KEY_REVIEW);
}
@Bean
public Binding docNotifyBinding() {
return BindingBuilder.bind(docNotifyQueue())
.to(docExchange())
.with(ROUTING_KEY_NOTIFY);
}
}
`
#### 3.3.2 消息发送示例
```java
@Service
@RequiredArgsConstructor
public class DocumentMessagePublisher {
private final RabbitTemplate rabbitTemplate;
public void sendParseMessage(DocumentMessage message) {
rabbitTemplate.convertAndSend(
RabbitMQConfig.DOC_EXCHANGE,
RabbitMQConfig.ROUTING_KEY_PARSE,
message,
msg -> {
msg.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
msg.getMessageProperties().setPriority(message.getPriority());
msg.getMessageProperties().setHeader("documentId", message.getDocumentId());
return msg;
}
);
}
public void sendReviewMessage(DocumentMessage message) {
rabbitTemplate.convertAndSend(
RabbitMQConfig.DOC_EXCHANGE,
RabbitMQConfig.ROUTING_KEY_REVIEW,
message
);
}
}
`
4. 文档处理任务调度

4.1 分布式任务调度概述
在分布式环境下,任务调度面临以下挑战:
1. **任务分配**:如何将任务均匀分配到各个工作节点
2. **状态管理**:如何跟踪任务执行状态
3. **故障处理**:节点故障时如何处理正在执行的任务
4. **资源竞争**:如何避免多个节点处理同一个任务
4.2 主流调度框架对比
|
框架 |
定时调度 |
任务分片 |
故障转移 |
管理界面 |
|
Xxl-Job |
支持 |
支持 |
支持 |
有 |
|
ElasticJob |
支持 |
支持 |
支持 |
有 |
|
Quartz |
支持 |
不支持 |
有限 |
无 |
本系统选择XXL-Job作为任务调度框架,原因如下:
- 社区活跃,文档完善
- 支持任务分片,适合大规模处理
- 内置管理界面,方便运维
- 轻量级,集成方便
4.3 XXL-Job集成配置
#### 4.3.1 调度中心配置
```yaml
application.yml
xxl:
job:
admin:
addresses: http://xxl-job-admin:8080/xxl-job-admin
executor:
appname: document-processor
port: 9999
logpath: /data/applogs/xxl-job/jobhandler
logretentiondays: 30
accessToken: default_token_for_security
`
#### 4.3.2 执行器配置
```java
@Configuration
public class XxlJobConfig {
@Bean
public XxlJobSpringExecutor xxlJobExecutor() {
XxlJobSpringExecutor xxlJobSpringExecutor = new XxlJobSpringExecutor();
xxlJobSpringExecutor.setAdminAddresses("http://xxl-job-admin:8080/xxl-job-admin");
xxlJobSpringExecutor.setAppname("document-processor");
xxlJobSpringExecutor.setPort(9999);
xxlJobSpringExecutor.setAccessToken("default_token_for_security");
xxlJobSpringExecutor.setLogPath("/data/applogs/xxl-job/jobhandler");
xxlJobSpringExecutor.setLogRetentionDays(30);
return xxlJobSpringExecutor;
}
}
`
4.4 任务分片处理
```java
@Component
public class DocumentShardJobHandler {
private final DocumentProcessingService documentProcessingService;
private final MongoTemplate mongoTemplate;
@XxlJob("documentShardJob")
public ReturnT<String> execute(String param) {
// 获取分片参数
ShardingUtil.ShardingVO shardingVO = ShardingUtil.getSharding();
int index = shardingVO.getIndex();
int total = shardingVO.getTotal();
log.info("开始执行文档分片任务,分片编号: {}/{}", index, total);
// 构建分片查询条件
Query query = new Query();
query.addCriteria(Criteria.where("status").is(DocumentStatus.PENDING));
query.addCriteria(Criteria.where("shardIndex").is(index));
query.limit(100); // 每批次处理100条
// 查询待处理文档
List<Document> documents = mongoTemplate.find(query, Document.class);
int successCount = 0;
int failCount = 0;
for (Document doc : documents) {
try {
documentProcessingService.processDocument(doc);
successCount++;
} catch (Exception e) {
log.error("处理文档失败: {}", doc.getId(), e);
failCount++;
}
}
log.info("文档分片任务完成,成功: {},失败: {}", successCount, failCount);
return new ReturnT<>(String.format("成功:%d, 失败:%d", successCount, failCount));
}
}
`
4.5 任务状态管理
任务状态机设计:
`
PENDING ──────► RUNNING ──────► SUCCEEDED
│
│
▼
FAILED ──────► RETRYING ──────► RUNNING
│
▼
CANCELLED
`
状态说明:
|
状态 |
含义 |
可转换状态 |
|
PENDING |
等待执行 |
RUNNING, CANCELLED |
|
RUNNING |
执行中 |
SUCCEEDED, FAILED, RETRYING |
|
SUCCEEDED |
执行成功 |
- |
|
FAILED |
执行失败 |
RETRYING, CANCELLED |
|
RETRYING |
重试中 |
RUNNING, CANCELLED |
|
CANCELLED |
已取消 |
- |
5. 失败重试与幂等处理

5.1 失败重试机制
#### 5.1.1 重试策略设计
系统在消息消费层实现智能重试机制:
```java
@Component
public class DocumentMessageConsumer {
private static final int MAX_RETRY_COUNT = 3;
private static final long INITIAL_BACKOFF_MS = 1000;
private static final double BACKOFF_MULTIPLIER = 2.0;
private final DocumentProcessingService documentProcessingService;
private final MessagePublisher messagePublisher;
@KafkaListener(topics = "doc-parse", groupId = "doc-processor")
public void consumeDocumentMessage(ConsumerRecord<String, DocumentMessage> record,
Acknowledgment acknowledgment) {
DocumentMessage message = record.value();
int retryCount = getRetryCount(record);
try {
log.info("开始处理文档消息,文档ID: {},重试次数: {}",
message.getDocumentId(), retryCount);
documentProcessingService.processDocument(message);
// 处理成功,确认消息
acknowledgment.acknowledge();
log.info("文档处理成功,文档ID: {}", message.getDocumentId());
} catch (Exception e) {
log.error("文档处理失败,文档ID: {},重试次数: {}",
message.getDocumentId(), retryCount, e);
if (retryCount < MAX_RETRY_COUNT) {
// 计算退避时间
long backoffMs = (long) (INITIAL_BACKOFF_MS *
Math.pow(BACKOFF_MULTIPLIER, retryCount));
// 延迟重发
scheduleRetry(message, retryCount + 1, backoffMs);
acknowledgment.acknowledge();
} else {
// 超过最大重试次数,发送到死信队列
messagePublisher.sendToDeadLetterQueue(message, e.getMessage());
acknowledgment.acknowledge();
}
}
}
private int getRetryCount(ConsumerRecord<String, DocumentMessage> record) {
Header retryHeader = record.headers().lastHeader("x-retry-count");
if (retryHeader != null) {
return Integer.parseInt(new String(retryHeader.value()));
}
return 0;
}
}
`
#### 5.1.2 退避策略
```java
public class BackoffStrategy {
/**
* 指数退避策略
delay = base multiplier^retryCount
*/
public static long exponentialBackoff(int retryCount, long base, double multiplier) {
return (long) (base * Math.pow(multiplier, retryCount));
}
/**
* 抖动策略,防止多节点同时重试
*/
public static long withJitter(long delay, double jitterFactor) {
Random random = new Random();
long jitter = (long) (delay jitterFactor random.nextDouble());
return delay + jitter;
}
/**
* 完整退避计算
*/
public static long calculateBackoff(int retryCount) {
long baseDelay = 1000; // 1秒
double multiplier = 2.0;
double jitterFactor = 0.2;
long exponentialDelay = exponentialBackoff(retryCount, baseDelay, multiplier);
return withJitter(exponentialDelay, jitterFactor);
}
}
`
5.2 幂等性保证
#### 5.2.1 幂等键设计
```java
@Component
public class IdempotencyService {
private final StringRedisTemplate redisTemplate;
private final IdempotencyRecordMapper idempotencyRecordMapper;
private static final String IDEMPOTENCY_KEY_PREFIX = "idempotency:";
private static final long IDEMPOTENCY_TTL_HOURS = 24;
/**
* 生成幂等键
*/
public String generateIdempotencyKey(String documentId, String operation) {
return String.format("%s:%s:%d", documentId, operation, System.currentTimeMillis() / 60000);
}
/**
* 检查并获取锁
*/
public boolean tryAcquire(String idempotencyKey, long expireSeconds) {
// 1. 尝试获取Redis分布式锁
Boolean acquired = redisTemplate.opsForValue()
.setIfAbsent(IDEMPOTENCY_KEY_PREFIX + idempotencyKey,
"PROCESSING",
Duration.ofSeconds(expireSeconds));
if (Boolean.TRUE.equals(acquired)) {
return true;
}
// 2. 检查是否为已完成的任务
String value = redisTemplate.opsForValue().get(IDEMPOTENCY_KEY_PREFIX + idempotencyKey);
return "COMPLETED".equals(value);
}
/**
* 标记处理完成
*/
public void markCompleted(String idempotencyKey) {
redisTemplate.opsForValue().set(
IDEMPOTENCY_KEY_PREFIX + idempotencyKey,
"COMPLETED",
Duration.ofHours(IDEMPOTENCY_TTL_HOURS)
);
}
/**
* 检查是否为重复请求
*/
public IdempotencyStatus checkStatus(String idempotencyKey) {
String value = redisTemplate.opsForValue().get(IDEMPOTENCY_KEY_PREFIX + idempotencyKey);
if (value == null) {
return IdempotencyStatus.NOT_FOUND;
}
return switch (value) {
case "PROCESSING" -> IdempotencyStatus.PROCESSING;
case "COMPLETED" -> IdempotencyStatus.COMPLETED;
case "FAILED" -> IdempotencyStatus.FAILED;
default -> IdempotencyStatus.UNKNOWN;
};
}
}
public enum IdempotencyStatus {
NOT_FOUND,
PROCESSING,
COMPLETED,
FAILED,
UNKNOWN
}
`
#### 5.2.2 消息处理幂等性实现
```java
@Service
@RequiredArgsConstructor
public class IdempotentDocumentProcessor {
private final IdempotencyService idempotencyService;
private final DocumentProcessingService documentProcessingService;
/**
* 幂等处理文档
*/
public ProcessingResult processIdempotent(DocumentMessage message) {
String idempotencyKey = buildIdempotencyKey(message);
// 1. 检查幂等状态
IdempotencyStatus status = idempotencyService.checkStatus(idempotencyKey);
if (status == IdempotencyStatus.COMPLETED) {
log.info("检测到重复消息,直接返回成功,文档ID: {}", message.getDocumentId());
return ProcessingResult.duplicate();
}
if (status == IdempotencyStatus.PROCESSING) {
log.info("检测到正在处理的请求,文档ID: {}", message.getDocumentId());
return ProcessingResult.processing();
}
// 2. 尝试获取处理权
if (!idempotencyService.tryAcquire(idempotencyKey, 300)) {
return ProcessingResult.busy();
}
try {
// 3. 执行实际处理
ProcessingResult result = documentProcessingService.processDocument(message);
// 4. 根据结果更新状态
if (result.isSuccess()) {
idempotencyService.markCompleted(idempotencyKey);
} else {
idempotencyService.markFailed(idempotencyKey);
}
return result;
} catch (Exception e) {
idempotencyService.markFailed(idempotencyKey);
throw e;
}
}
private String buildIdempotencyKey(DocumentMessage message) {
return String.format("doc:%s:op:%s",
message.getDocumentId(),
message.getOperationType());
}
}
`
6. 实际代码示例:消息生产者/消费者实现
6.1 消息生产者完整实现
```java
/**
* 文档消息生产者
* 负责将文档处理任务发送到消息队列
*/
@Service
@Slf4j
public class DocumentMessageProducer {
private final KafkaTemplate<String, DocumentMessage> kafkaTemplate;
private final RabbitTemplate rabbitTemplate;
// Kafka Topic常量
private static final String KAFKA_TOPIC_DOC_PARSE = "doc-parse";
private static final String KAFKA_TOPIC_DOC_REVIEW = "doc-review";
private static final String KAFKA_TOPIC_DOC_NOTIFY = "doc-notify";
// RabbitMQ常量
private static final String DOC_EXCHANGE = "doc.exchange";
private static final String ROUTING_KEY_PARSE = "doc.parse";
private static final String ROUTING_KEY_REVIEW = "doc.review";
/**
* 发送文档解析任务(Kafka)
*/
public void sendParseTask(DocumentMessage message) {
log.info("发送文档解析任务,文档ID: {}, 优先级: {}",
message.getDocumentId(), message.getPriority());
// 设置发送时间戳
message.setTimestamp(System.currentTimeMillis());
// 使用文档ID作为key,保证同一文档的消息发送到同一分区
kafkaTemplate.send(KAFKA_TOPIC_DOC_PARSE,
message.getDocumentId(),
message)
.whenComplete((result, ex) -> {
if (ex != null) {
log.error("发送文档解析任务失败,文档ID: {}",
message.getDocumentId(), ex);
} else {
log.info("发送文档解析任务成功,分区: {}, 偏移量: {}",
result.getRecordMetadata().partition(),
result.getRecordMetadata().offset());
}
});
}
/**
* 发送文档审核任务(Kafka)
*/
public void sendReviewTask(DocumentMessage message) {
log.info("发送文档审核任务,文档ID: {}", message.getDocumentId());
message.setTimestamp(System.currentTimeMillis());
kafkaTemplate.send(KAFKA_TOPIC_DOC_REVIEW,
message.getDocumentId(),
message);
}
/**
* 发送通知任务(RabbitMQ)
*/
public void sendNotifyTask(DocumentMessage message) {
log.info("发送通知任务,文档ID: {}, 通知类型: {}",
message.getDocumentId(), message.getNotifyType());
rabbitTemplate.convertAndSend(DOC_EXCHANGE,
ROUTING_KEY_PARSE,
message,
msg -> {
msg.getMessageProperties()
.setDeliveryMode(MessageDeliveryMode.PERSISTENT);
msg.getMessageProperties()
.setPriority(message.getPriority());
return msg;
});
}
/**
* 发送批量任务
*/
public void sendBatchTasks(List<DocumentMessage> messages) {
log.info("发送批量任务,数量: {}", messages.size());
List<Future<SendResult<String, DocumentMessage>>> futures = messages.stream()
.map(msg -> {
msg.setTimestamp(System.currentTimeMillis());
return kafkaTemplate.send(KAFKA_TOPIC_DOC_PARSE,
msg.getDocumentId(),
msg);
})
.collect(Collectors.toList());
// 等待所有发送完成
for (int i = 0; i < futures.size(); i++) {
try {
futures.get(i).get(10, TimeUnit.SECONDS);
} catch (Exception e) {
log.error("批量发送任务失败,索引: {}", i, e);
}
}
}
}
`
6.2 消息消费者完整实现
```java
/**
* 文档消息消费者
* 负责从消息队列消费任务并执行处理
*/
@Service
@Slf4j
public class DocumentMessageConsumer {
private final DocumentProcessingService processingService;
private final IdempotencyService idempotencyService;
private final DeadLetterHandler deadLetterHandler;
// 并发配置
private static final int CONCURRENT_LISTENERS = 3;
private static final int MAX_POLL_RECORDS = 50;
/**
* 消费文档解析消息
*/
@KafkaListener(
topics = "doc-parse",
groupId = "doc-processor-group",
containerFactory = "kafkaListenerContainerFactory"
)
public void consumeParseMessage(ConsumerRecord<String, DocumentMessage> record,
Acknowledgment acknowledgment) {
DocumentMessage message = record.value();
log.info("消费文档解析消息,文档ID: {}, 分区: {}, 偏移量: {}",
message.getDocumentId(), record.partition(), record.offset());
try {
// 幂等处理
ProcessingResult result = processWithIdempotency(message);
if (result.isSuccess()) {
log.info("文档解析成功,文档ID: {}, 处理结果: {}",
message.getDocumentId(), result.getMessage());
} else if (result.isDuplicate()) {
log.info("检测到重复消息,跳过处理,文档ID: {}", message.getDocumentId());
} else {
log.warn("文档解析处理返回异常结果,文档ID: {}, 结果: {}",
message.getDocumentId(), result.getMessage());
}
acknowledgment.acknowledge();
} catch (Exception e) {
log.error("消费文档解析消息异常,文档ID: {}", message.getDocumentId(), e);
handleConsumeException(message, e, record, acknowledgment);
}
}
/**
* 消费文档审核消息
*/
@KafkaListener(
topics = "doc-review",
groupId = "doc-reviewer-group",
containerFactory = "kafkaListenerContainerFactory"
)
public void consumeReviewMessage(ConsumerRecord<String, DocumentMessage> record,
Acknowledgment acknowledgment) {
DocumentMessage message = record.value();
try {
processingService.reviewDocument(message);
acknowledgment.acknowledge();
} catch (Exception e) {
log.error("消费文档审核消息异常,文档ID: {}", message.getDocumentId(), e);
handleConsumeException(message, e, record, acknowledgment);
}
}
/**
* 消费通知消息(RabbitMQ)
*/
@RabbitListener(queues = "doc.notify.queue")
public void consumeNotifyMessage(DocumentMessage message,
Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) {
log.info("消费通知消息,文档ID: {}, 通知类型: {}",
message.getDocumentId(), message.getNotifyType());
try {
processingService.sendNotification(message);
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
log.error("消费通知消息异常,文档ID: {}", message.getDocumentId(), e);
// 通知消息通常不重试,直接拒绝
channel.basicNack(deliveryTag, false, false);
}
}
/**
* 带幂等性的消息处理
*/
private ProcessingResult processWithIdempotency(DocumentMessage message) {
String idempotencyKey = buildIdempotencyKey(message);
// 检查是否已处理
IdempotencyStatus status = idempotencyService.checkStatus(idempotencyKey);
if (status == IdempotencyStatus.COMPLETED) {
return ProcessingResult.duplicate();
}
// 尝试获取处理权
if (!idempotencyService.tryAcquire(idempotencyKey, 300)) {
return ProcessingResult.busy();
}
try {
ProcessingResult result = processingService.processDocument(message);
if (result.isSuccess()) {
idempotencyService.markCompleted(idempotencyKey);
}
return result;
} catch (Exception e) {
idempotencyService.markFailed(idempotencyKey);
throw e;
}
}
/**
* 处理消费异常
*/
private void handleConsumeException(DocumentMessage message,
Exception e,
ConsumerRecord<String, DocumentMessage> record,
Acknowledgment acknowledgment) {
int retryCount = getRetryCount(record);
if (retryCount < 3) {
// 重新入队重试
log.warn("消息处理失败,准备重试,当前重试次数: {}", retryCount);
// 实际重试逻辑由重试框架处理
acknowledgment.acknowledge();
} else {
// 进入死信队列
log.error("消息处理失败次数过多,发送到死信队列,文档ID: {}", message.getDocumentId());
deadLetterHandler.sendToDeadLetterQueue(message, e.getMessage());
acknowledgment.acknowledge();
}
}
private String buildIdempotencyKey(DocumentMessage message) {
return String.format("doc:%s:partition:%d:offset:%d",
message.getDocumentId(),
message.getPartition(),
message.getOffset());
}
private int getRetryCount(ConsumerRecord<String, DocumentMessage> record) {
Header retryHeader = record.headers().lastHeader("x-retry-count");
if (retryHeader != null) {
return Integer.parseInt(new String(retryHeader.value()));
}
return 0;
}
}
`
6.3 消息DTO定义
```java
/**
* 文档消息定义
*/
@Data
@Builder
public class DocumentMessage implements Serializable {
private static final long serialVersionUID = 1L;
/**
* 文档ID
*/
private String documentId;
/**
* 文档名称
*/
private String documentName;
/**
* 文档类型
*/
private DocumentType documentType;
/**
* 操作类型
*/
private OperationType operationType;
/**
* 消息优先级 0-9,数字越大优先级越高
*/
private int priority;
/**
* 通知类型
*/
private NotifyType notifyType;
/**
* 消息时间戳
*/
private long timestamp;
/**
* 分区编号
*/
private int partition;
/**
* 消息偏移量
*/
private long offset;
/**
* 重试次数
*/
@Builder.Default
private int retryCount = 0;
/**
* 扩展属性
*/
private Map<String, String> extProps;
}
public enum OperationType {
PARSE, // 解析
REVIEW, // 审核
NOTIFY, // 通知
ARCHIVE // 归档
}
public enum NotifyType {
EMAIL,
SMS,
WEBHOOK,
IN_APP
}
`
7. 实际代码示例:任务调度服务
7.1 任务调度服务实现
```java
/**
* 文档处理任务调度服务
*/
@Service
@Slf4j
public class DocumentTaskScheduler {
private final DocumentTaskMapper taskMapper;
private final DocumentMessageProducer messageProducer;
private final RedisTemplate<String, String> redisTemplate;
private static final String TASK_LOCK_KEY = "doc:scheduler:lock";
private static final long LOCK_EXPIRE_SECONDS = 300;
/**
* 调度待处理任务
* 每分钟执行一次,将待处理任务发送到消息队列
*/
@XxlJob("schedulePendingTasksJob")
public ReturnT<String> schedulePendingTasks() {
log.info("开始执行任务调度...");
// 1. 获取分布式锁,防止多节点重复调度
Boolean lockAcquired = redisTemplate.opsForValue()
.setIfAbsent(TASK_LOCK_KEY,
Thread.currentThread().getName(),
Duration.ofSeconds(LOCK_EXPIRE_SECONDS));
if (!Boolean.TRUE.equals(lockAcquired)) {
log.info("未获取到调度锁,跳过本次调度");
return new ReturnT<>(ReturnT.SUCCESS_CODE, "skip");
}
try {
// 2. 查询待调度任务
List<DocumentTask> pendingTasks = taskMapper.selectPendingTasks(1000);
log.info("查询到待调度任务数量: {}", pendingTasks.size());
int scheduledCount = 0;
int failedCount = 0;
for (DocumentTask task : pendingTasks) {
try {
// 3. 更新任务状态为调度中
task.setStatus(TaskStatus.SCHEDULING);
task.setScheduleTime(System.currentTimeMillis());
taskMapper.updateById(task);
// 4. 发送消息到队列
DocumentMessage message = buildMessage(task);
if (task.getTaskType() == TaskType.PARSE) {
messageProducer.sendParseTask(message);
} else if (task.getTaskType() == TaskType.REVIEW) {
messageProducer.sendReviewTask(message);
}
// 5. 更新任务状态为已调度
task.setStatus(TaskStatus.SCHEDULED);
task.setMessageId(message.getMessageId());
taskMapper.updateById(task);
scheduledCount++;
} catch (Exception e) {
log.error("调度任务失败,任务ID: {}", task.getId(), e);
task.setStatus(TaskStatus.FAILED);
task.setErrorMessage(e.getMessage());
taskMapper.updateById(task);
failedCount++;
}
}
log.info("任务调度完成,成功: {},失败: {}", scheduledCount, failedCount);
return new ReturnT<>(String.format("成功:%d, 失败:%d", scheduledCount, failedCount));
} finally {
// 释放分布式锁
redisTemplate.delete(TASK_LOCK_KEY);
}
}
/**
* 清理超时任务
* 任务调度后超过指定时间未完成,标记为超时
*/
@XxlJob("cleanupTimeoutTasksJob")
public ReturnT<String> cleanupTimeoutTasks() {
log.info("开始执行超时任务清理...");
long timeoutMillis = 30 60 1000; // 30分钟
long cutoffTime = System.currentTimeMillis() - timeoutMillis;
List<DocumentTask> timeoutTasks = taskMapper.selectTimeoutTasks(cutoffTime);
log.info("查询到超时任务数量: {}", timeoutTasks.size());
for (DocumentTask task : timeoutTasks) {
task.setStatus(TaskStatus.TIMEOUT);
task.setErrorMessage("任务执行超时");
taskMapper.updateById(task);
// 发送告警通知
sendTimeoutAlert(task);
}
return new ReturnT<>(String.format("清理超时任务:%d", timeoutTasks.size()));
}
/**
* 重试失败任务
*/
@XxlJob("retryFailedTasksJob")
public ReturnT<String> retryFailedTasks() {
log.info("开始执行失败任务重试...");
List<DocumentTask> failedTasks = taskMapper.selectFailedTasks(100);
log.info("查询到失败任务数量: {}", failedTasks.size());
int retriedCount = 0;
for (DocumentTask task : failedTasks) {
// 检查重试次数
if (task.getRetryCount() >= 3) {
log.warn("任务重试次数已达上限,任务ID: {}", task.getId());
task.setStatus(TaskStatus.EXHAUSTED);
taskMapper.updateById(task);
continue;
}
// 重置状态并增加重试次数
task.setStatus(TaskStatus.PENDING);
task.setRetryCount(task.getRetryCount() + 1);
task.setLastRetryTime(System.currentTimeMillis());
taskMapper.updateById(task);
// 重新发送消息
DocumentMessage message = buildMessage(task);
message.setRetryCount(task.getRetryCount());
if (task.getTaskType() == TaskType.PARSE) {
messageProducer.sendParseTask(message);
}
retriedCount++;
}
return new ReturnT<>(String.format("重试任务:%d", retriedCount));
}
private DocumentMessage buildMessage(DocumentTask task) {
return DocumentMessage.builder()
.documentId(task.getDocumentId())
.documentName(task.getDocumentName())
.documentType(task.getDocumentType())
.operationType(task.getOperationType())
.priority(task.getPriority())
.timestamp(System.currentTimeMillis())
.retryCount(task.getRetryCount())
.extProps(task.getExtProps())
.build();
}
private void sendTimeoutAlert(DocumentTask task) {
log.warn("任务执行超时告警,任务ID: {}, 文档ID: {}, 执行时间: {}ms",
task.getId(), task.getDocumentId(),
System.currentTimeMillis() - task.getScheduleTime());
}
}
`
7.2 任务状态监听器
```java
/**
* 任务状态监听器
* 监听任务状态变化,执行相应操作
*/
@Component
@Slf4j
public class TaskStatusListener {
private final ApplicationEventPublisher eventPublisher;
private final NotificationService notificationService;
@EventListener
public void handleTaskStatusChange(TaskStatusChangeEvent event) {
DocumentTask task = event.getTask();
TaskStatus oldStatus = event.getOldStatus();
TaskStatus newStatus = event.getNewStatus();
log.info("任务状态变更,任务ID: {}, {} -> {}",
task.getId(), oldStatus, newStatus);
switch (newStatus) {
case SUCCEEDED -> onTaskSucceeded(task);
case FAILED -> onTaskFailed(task);
case TIMEOUT -> onTaskTimeout(task);
case CANCELLED -> onTaskCancelled(task);
default -> {}
}
}
private void onTaskSucceeded(DocumentTask task) {
log.info("任务执行成功,任务ID: {}, 文档ID: {}",
task.getId(), task.getDocumentId());
// 发送成功通知
notificationService.sendTaskNotification(
NotificationType.TASK_SUCCESS,
task.getDocumentId(),
"任务执行成功"
);
}
private void onTaskFailed(DocumentTask task) {
log.error("任务执行失败,任务ID: {}, 错误信息: {}",
task.getId(), task.getErrorMessage());
// 发送失败告警
notificationService.sendTaskNotification(
NotificationType.TASK_FAILED,
task.getDocumentId(),
task.getErrorMessage()
);
}
private void onTaskTimeout(DocumentTask task) {
log.warn("任务执行超时,任务ID: {}, 文档ID: {}",
task.getId(), task.getDocumentId());
// 发送超时告警
notificationService.sendTaskNotification(
NotificationType.TASK_TIMEOUT,
task.getDocumentId(),
"任务执行超时"
);
}
private void onTaskCancelled(DocumentTask task) {
log.info("任务被取消,任务ID: {}", task.getId());
}
}
`
8. 实际代码示例:分布式锁实现
8.1 Redis分布式锁
```java
/**
* Redis分布式锁实现
*/
@Component
@RequiredArgsConstructor
public class RedisDistributedLock {
private final StringRedisTemplate redisTemplate;
private static final String LOCK_PREFIX = "lock:";
private static final String LOCK_VALUE_PREFIX = Thread.currentThread().getName() + ":" + UUID.randomUUID().toString();
/**
* 尝试获取锁
*/
public boolean tryLock(String key, long expireSeconds) {
String lockKey = LOCK_PREFIX + key;
String lockValue = LOCK_VALUE_PREFIX + ":" + System.currentTimeMillis();
Boolean result = redisTemplate.opsForValue()
.setIfAbsent(lockKey, lockValue, Duration.ofSeconds(expireSeconds));
return Boolean.TRUE.equals(result);
}
/**
* 尝试获取锁(带重试)
*/
public boolean tryLockWithRetry(String key, long expireSeconds,
int maxRetries, long retryIntervalMs) {
for (int i = 0; i < maxRetries; i++) {
if (tryLock(key, expireSeconds)) {
return true;
}
try {
Thread.sleep(retryIntervalMs);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return false;
}
}
return false;
}
/**
* 释放锁
*/
public boolean unlock(String key) {
String lockKey = LOCK_PREFIX + key;
return redisTemplate.delete(lockKey);
}
/**
* 释放锁(安全释放,只释放自己的锁)
*/
public boolean safeUnlock(String key, String expectedValue) {
String lockKey = LOCK_PREFIX + key;
String currentValue = redisTemplate.opsForValue().get(lockKey);
if (expectedValue.equals(currentValue)) {
return redisTemplate.delete(lockKey);
}
return false;
}
/**
* 延长锁的过期时间
*/
public boolean extendLock(String key, long additionalSeconds) {
String lockKey = LOCK_PREFIX + key;
return Boolean.TRUE.equals(
redisTemplate.expire(lockKey, Duration.ofSeconds(additionalSeconds))
);
}
/**
* 执行带锁的操作
*/
public <T> T executeWithLock(String key, long expireSeconds,
Callable<T> action) throws Exception {
if (!tryLock(key, expireSeconds)) {
throw new LockAcquisitionException("无法获取锁: " + key);
}
try {
return action.call();
} finally {
unlock(key);
}
}
}
`
8.2 数据库分布式锁
```java
/**
* 基于数据库的分布式锁实现
*/
@Component
@RequiredArgsConstructor
public class DatabaseDistributedLock {
private final JdbcTemplate jdbcTemplate;
private static final String LOCK_TABLE = "distributed_locks";
/**
* 尝试获取锁
*/
public boolean tryLock(String lockName, String lockOwner, long expireSeconds) {
String sql = """
INSERT INTO distributed_locks (lock_name, lock_owner, expire_time, created_at)
VALUES (?, ?, ?, NOW())
ON DUPLICATE KEY UPDATE
lock_owner = IF(expire_time < NOW(), VALUES(lock_owner), lock_owner),
expire_time = IF(lock_owner = VALUES(lock_owner) OR expire_time < NOW(),
VALUES(expire_time), expire_time)
""";
// 检查锁是否获取成功
String checkSql = """
SELECT COUNT(*) FROM distributed_locks
WHERE lock_name = ? AND lock_owner = ?
AND expire_time > NOW()
""";
jdbcTemplate.update(sql, lockName, lockOwner);
Integer count = jdbcTemplate.queryForObject(checkSql, Integer.class,
lockName, lockOwner);
return count != null && count > 0;
}
/**
* 释放锁
*/
public boolean unlock(String lockName, String lockOwner) {
String sql = """
DELETE FROM distributed_locks
WHERE lock_name = ? AND lock_owner = ?
""";
int affectedRows = jdbcTemplate.update(sql, lockName, lockOwner);
return affectedRows > 0;
}
/**
* 检查锁状态
*/
public LockStatus checkLockStatus(String lockName) {
String sql = """
SELECT lock_owner, expire_time FROM distributed_locks
WHERE lock_name = ?
""";
List<LockInfo> locks = jdbcTemplate.query(sql,
(rs, rowNum) -> new LockInfo(
rs.getString("lock_owner"),
rs.getTimestamp("expire_time").toInstant().toEpochMilli()
),
lockName
);
if (locks.isEmpty()) {
return LockStatus.NOT_EXISTS;
}
LockInfo lock = locks.get(0);
if (lock.getExpireTime() < System.currentTimeMillis()) {
return LockStatus.EXPIRED;
}
return LockStatus.HELD;
}
/**
* 清理过期锁
*/
public int cleanupExpiredLocks() {
String sql = "DELETE FROM distributed_locks WHERE expire_time < NOW()";
return jdbcTemplate.update(sql);
}
@Data
@AllArgsConstructor
private static class LockInfo {
private String owner;
private long expireTime;
}
public enum LockStatus {
NOT_EXISTS,
HELD,
EXPIRED
}
}
`
8.3 Zookeeper分布式锁
```java
/**
* 基于Zookeeper的分布式锁实现
*/
@Component
@Slf4j
public class ZookeeperDistributedLock {
private final CuratorFramework curatorFramework;
private static final String LOCK_PATH_PREFIX = "/locks/";
/**
* 创建临时顺序节点获取锁
*/
public String tryAcquireLock(String lockName, long timeoutMs) throws Exception {
String lockPath = LOCK_PATH_PREFIX + lockName;
// 确保父节点存在
if (curatorFramework.checkExists().forPath(lockPath) == null) {
curatorFramework.create()
.creatingParentsIfNeeded()
.withMode(CreateMode.PERSISTENT)
.forPath(lockPath);
}
// 创建临时顺序节点
String nodePath = curatorFramework.create()
.withMode(CreateMode.EPHEMERAL_SEQUENTIAL)
.forPath(lockPath + "/lock-");
// 检查是否为最小节点
if (isLowestNode(nodePath)) {
log.debug("成功获取锁: {}, 节点: {}", lockName, nodePath);
return nodePath;
}
// 等待前面的节点释放
awaitPreviousNode(nodePath, timeoutMs);
return nodePath;
}
/**
* 释放锁
*/
public void releaseLock(String nodePath) {
try {
curatorFramework.delete().forPath(nodePath);
log.debug("释放锁成功: {}", nodePath);
} catch (Exception e) {
log.error("释放锁失败: {}", nodePath, e);
}
}
/**
* 检查是否为最小节点
*/
private boolean isLowestNode(String nodePath) throws Exception {
String parentPath = nodePath.substring(0, nodePath.lastIndexOf('/'));
List<String> children = curatorFramework.getChildren().forPath(parentPath);
Collections.sort(children);
String smallestNode = children.get(0);
return nodePath.endsWith(smallestNode);
}
/**
* 等待前面的节点释放
*/
private void awaitPreviousNode(String nodePath, long timeoutMs) throws Exception {
String parentPath = nodePath.substring(0, nodePath.lastIndexOf('/'));
while (true) {
List<String> children = curatorFramework.getChildren().forPath(parentPath);
Collections.sort(children);
String currentNode = nodePath.substring(nodePath.lastIndexOf('/') + 1);
int currentIndex = children.indexOf(currentNode);
if (currentIndex <= 0) {
// 当前节点已经是最小的
return;
}
// 监听前一个节点的删除事件
String previousNode = children.get(currentIndex - 1);
String previousPath = parentPath + "/" + previousNode;
CountDownLatch latch = new CountDownLatch(1);
curatorFramework.getData().usingWatcher(
event -> {
if (event.getType() == Watcher.Event.EventType.NodeDeleted) {
latch.countDown();
}
}
).forPath(previousPath);
if (latch.await(timeoutMs, TimeUnit.MILLISECONDS)) {
return;
}
throw new TimeoutException("等待锁超时");
}
}
}
`
9. 章节总结
9.1 核心要点回顾
本章节详细介绍了文档智能解析审核系统的分布式架构设计,主要内容包括:
1. **分布式架构设计**
- 采用微服务架构,实现服务拆分和独立扩展
- 通过API网关实现统一入口和负载均衡
- 多元存储策略满足不同数据存储需求
2. **消息队列集成**
- Kafka适用于高吞吐量场景
- RabbitMQ适用于需要灵活路由的场景
- 详细配置了生产者和消费者
3. **任务调度系统**
- 使用XXL-Job实现分布式任务调度
- 支持任务分片处理
- 完善的任务状态管理
4. **失败重试与幂等处理**
- 指数退避重试策略
- 多级缓存实现幂等性保证
- 死信队列处理异常消息
9.2 进阶学习建议
- 深入学习Kubernetes容器编排,实现自动扩缩容
- 研究服务网格(Service Mesh)如Istio,提升服务治理能力
- 学习事件驱动架构(EDA),进一步提升系统解耦程度
版权声明:本文为洛水石原创技术文章,版权所有,未经许可禁止转载。
作者:洛水石 | 文档智能解析审核
AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。
更多推荐



所有评论(0)