本文目录导读:

在Java中实现动态发布案例,我为你提供一个完整的实现方案,涵盖常见的几种动态发布场景。
完整的动态发布系统实现
// 1. 内容实体类
package com.example.publish.entity;
import lombok.Data;
import java.time.LocalDateTime;
import java.util.List;
@Data
public class Content {
private String id;
private String title;
private String content;
private ContentType type;
private ContentStatus status;
private String author;
private List<String> tags;
private LocalDateTime publishTime;
private LocalDateTime createTime;
private LocalDateTime updateTime;
private Integer version;
private Integer channelId;
public enum ContentType {
ARTICLE, VIDEO, IMAGE, AUDIO, TEXT
}
public enum ContentStatus {
DRAFT, // 草稿
PENDING, // 待审核
APPROVED, // 已审核
PUBLISHED, // 已发布
OFFLINE, // 已下线
REJECTED // 已驳回
}
}
动态发布核心接口
// 2. 动态发布接口
package com.example.publish.service;
import com.example.publish.entity.Content;
import java.util.List;
public interface DynamicPublishService {
// 保存草稿
Content saveDraft(Content content);
// 提交审核
Content submitForReview(String contentId);
// 审核通过
Content approve(String contentId, String reviewer);
// 审核驳回
Content reject(String contentId, String reviewer, String reason);
// 动态发布
Content publish(String contentId);
// 定时发布
Content schedulePublish(String contentId, LocalDateTime publishTime);
// 下线内容
Content offline(String contentId);
// 更新内容
Content updateContent(Content content);
// 按状态查询
List<Content> queryByStatus(Content.ContentStatus status);
// 按渠道发布
Content publishToChannel(String contentId, String channelId);
}
动态发布实现类
// 3. 动态发布实现类
package com.example.publish.service.impl;
import com.example.publish.entity.Content;
import com.example.publish.service.DynamicPublishService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import java.time.LocalDateTime;
import java.util.List;
import java.util.UUID;
@Service
public class DynamicPublishServiceImpl implements DynamicPublishService {
@Autowired
private ContentMapper contentMapper;
@Autowired
private PublishChannelManager channelManager;
@Autowired
private RedisTemplate<String, Object> redisTemplate;
@Override
@Transactional
public Content saveDraft(Content content) {
if (content.getId() == null) {
content.setId(UUID.randomUUID().toString());
content.setStatus(Content.ContentStatus.DRAFT);
content.setCreateTime(LocalDateTime.now());
}
content.setUpdateTime(LocalDateTime.now());
content.setVersion(content.getVersion() == null ? 1 : content.getVersion() + 1);
contentMapper.save(content);
// 保存到Redis缓存
cacheContent(content);
return content;
}
@Override
@Transactional
public Content submitForReview(String contentId) {
Content content = getContent(contentId);
if (content.getStatus() != Content.ContentStatus.DRAFT) {
throw new IllegalStateException("只有草稿状态才能提交审核");
}
content.setStatus(Content.ContentStatus.PENDING);
content.setUpdateTime(LocalDateTime.now());
contentMapper.updateStatus(content);
// 通知审核人员
notifyReviewer(content);
return content;
}
@Override
public Content publish(String contentId) {
Content content = getContent(contentId);
// 检查状态是否是已审核或定时发布到期
if (content.getStatus() != Content.ContentStatus.APPROVED &&
content.getStatus() != Content.ContentStatus.PENDING) {
throw new IllegalStateException("内容状态不允许发布");
}
// 更新状态为已发布
content.setStatus(Content.ContentStatus.PUBLISHED);
content.setPublishTime(LocalDateTime.now());
content.setUpdateTime(LocalDateTime.now());
contentMapper.update(content);
// 推送到各种渠道
publishToChannels(content);
// 更新缓存
updateContentCache(content);
// 发送消息通知
sendPublishEvent(content);
return content;
}
@Override
public Content schedulePublish(String contentId, LocalDateTime publishTime) {
Content content = getContent(contentId);
content.setPublishTime(publishTime);
content.setStatus(Content.ContentStatus.APPROVED);
contentMapper.update(content);
// 创建定时任务
schedulePublishTask(contentId, publishTime);
return content;
}
@Override
public Content offline(String contentId) {
Content content = getContent(contentId);
content.setStatus(Content.ContentStatus.OFFLINE);
content.setUpdateTime(LocalDateTime.now());
contentMapper.update(content);
// 清理缓存
clearContentCache(contentId);
// 通知相关渠道下线内容
offlineFromChannels(content);
return content;
}
@Override
@Transactional
public Content updateContent(Content newContent) {
Content oldContent = getContent(newContent.getId());
// 保存历史版本
saveVersionHistory(oldContent);
// 更新内容
oldContent.setTitle(newContent.getTitle());
oldContent.setContent(newContent.getContent());
oldContent.setTags(newContent.getTags());
oldContent.setUpdateTime(LocalDateTime.now());
oldContent.setVersion(oldContent.getVersion() + 1);
contentMapper.update(oldContent);
// 更新缓存
cacheContent(oldContent);
return oldContent;
}
// 私有方法:获取内容
private Content getContent(String contentId) {
// 优先从缓存获取
Content content = (Content) redisTemplate.opsForValue()
.get("content:" + contentId);
if (content == null) {
content = contentMapper.findById(contentId);
if (content != null) {
cacheContent(content);
}
}
return content;
}
// 私有方法:缓存内容
private void cacheContent(Content content) {
redisTemplate.opsForValue().set(
"content:" + content.getId(),
content,
30, TimeUnit.MINUTES
);
}
// 私有方法:发布到多个渠道
private void publishToChannels(Content content) {
List<PublishChannel> channels =
channelManager.getEnabledChannels();
for (PublishChannel channel : channels) {
try {
channel.publish(content);
} catch (Exception e) {
// 记录发布失败
log.error("Publish to channel {} failed",
channel.getName(), e);
// 记录发布日志
savePublishLog(content.getId(),
channel.getName(), false, e.getMessage());
}
}
}
// 创建定时发布任务
private void schedulePublishTask(String contentId, LocalDateTime publishTime) {
// 使用Spring Task或Quartz
ScheduledExecutorService scheduler =
Executors.newScheduledThreadPool(4);
long delay = Duration.between(LocalDateTime.now(), publishTime).getSeconds();
scheduler.schedule(() -> {
publish(contentId);
}, delay, TimeUnit.SECONDS);
}
}
发布渠道接口和实现
// 4. 发布渠道接口
package com.example.publish.channel;
import com.example.publish.entity.Content;
public interface PublishChannel {
String getName();
boolean isEnabled();
void publish(Content content) throws Exception;
void offline(Content content) throws Exception;
}
// 5. 具体渠道实现
package com.example.publish.channel;
import com.example.publish.entity.Content;
import org.springframework.stereotype.Component;
@Component
public class SmsPublishChannel implements PublishChannel {
@Override
public String getName() {
return "SMS";
}
@Override
public boolean isEnabled() {
return config.isSmsEnabled();
}
@Override
public void publish(Content content) {
// 实现短信发送逻辑
sendSms(content.getContent());
}
@Override
public void offline(Content content) {
// 短信没有下线逻辑,可以留空
}
}
@Component
public class WebPublishChannel implements PublishChannel {
// Web端的发布逻辑
}
@Component
public class AppPublishChannel implements PublishChannel {
// App推送的发布逻辑
}
定时任务调度器
// 6. 定时发布任务调度器
package com.example.publish.scheduler;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
@Component
public class PublishScheduler {
@Autowired
private DynamicPublishService publishService;
@Autowired
private ScheduledContentRepository repository;
// 每30秒检查一次定时发布任务
@Scheduled(cron = "0/30 * * * * ?")
public void checkScheduledPublish() {
LocalDateTime now = LocalDateTime.now();
List<ScheduledContent> dueContentList =
repository.findByPublishTimeBeforeAndExecutedFalse(now);
for (ScheduledContent scheduled : dueContentList) {
try {
// 执行发布
publishService.publish(scheduled.getContentId());
// 标记为已执行
scheduled.setExecuted(true);
repository.save(scheduled);
} catch (Exception e) {
log.error("Scheduled publish failed for content: {}",
scheduled.getContentId(), e);
// 记录失败并可能需要重试
scheduled.setFailedCount(scheduled.getFailedCount() + 1);
repository.save(scheduled);
}
}
}
}
消息队列异步处理
// 7. 发布消息生产者
package com.example.publish.mq;
@Component
public class PublishMessageProducer {
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
public void sendPublishMessage(Content content) {
// 发送到Kafka消息队列
PublishMessage message = new PublishMessage();
message.setContentId(content.getId());
message.setPublishTime(content.getPublishTime());
kafkaTemplate.send("publish-topic",
content.getId(), message);
}
}
// 8. 发布消息消费者
@Component
public class PublishMessageConsumer {
@Autowired
private DynamicPublishService publishService;
@KafkaListener(topics = "publish-topic",
groupId = "publish-group")
public void handlePublishMessage(PublishMessage message) {
try {
// 执行发布
publishService.publish(message.getContentId());
} catch (Exception e) {
// 失败重试或记录日志
log.error("Failed to publish content from message", e);
// 可以发送到死信队列
}
}
}
使用示例
// 9. 使用示例
package com.example.publish.controller;
@RestController
@RequestMapping("/api/publish")
public class PublishController {
@Autowired
private DynamicPublishService publishService;
// 保存草稿
@PostMapping("/draft")
public Content saveDraft(@RequestBody Content content) {
return publishService.saveDraft(content);
}
// 提交审核
@PostMapping("/{id}/submit")
public Content submit(@PathVariable String id) {
return publishService.submitForReview(id);
}
// 立即发布
@PostMapping("/{id}/publish")
public Content publish(@PathVariable String id) {
return publishService.publish(id);
}
// 定时发布
@PostMapping("/{id}/schedule")
public Content schedulePublish(
@PathVariable String id,
@RequestParam LocalDateTime publishTime) {
return publishService.schedulePublish(id, publishTime);
}
// 审核
@PostMapping("/{id}/approve")
public Content approve(@PathVariable String id) {
return publishService.approve(id, "reviewer1");
}
// 下线
@PostMapping("/{id}/offline")
public Content offline(@PathVariable String id) {
return publishService.offline(id);
}
}
配置文件
# application.yml 配置
spring:
redis:
host: localhost
port: 6379
kafka:
bootstrap-servers: localhost:9092
# 发布渠道配置
publish:
channels:
sms:
enabled: true
gateway: "sms-gateway"
web:
enabled: true
cdn: "cdn-web"
app:
enabled: true
schedule:
pool-size: 4
check-interval: 30000 # 毫秒
这个完整的动态发布系统实现了: 管理支持草稿、审核、发布、下线的完整生命周期 2. 多渠道发布支持短信、Web、App等不同渠道 3. 定时发布支持精确到秒的定时发布 4. 消息队列使用Kafka进行异步处理,保证系统解耦 5. 缓存机制使用Redis缓存热点数据 6. 版本管理保存历史版本,支持回滚 7. 失败重试**:处理各种发布失败情况
你可以根据实际业务需求调整代码,增加更多功能如权限控制、统计分析等。