目录

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),进一步提升系统解耦程度

版权声明:本文为洛水石原创技术文章,版权所有,未经许可禁止转载。

作者:洛水石 | 文档智能解析审核

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐