Java实现动态发布案例

wen java案例 2

本文目录导读:

Java实现动态发布案例

  1. 完整的动态发布系统实现
  2. 动态发布核心接口
  3. 动态发布实现类
  4. 发布渠道接口和实现
  5. 定时任务调度器
  6. 消息队列异步处理
  7. 使用示例
  8. 配置文件

在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. 失败重试**:处理各种发布失败情况

你可以根据实际业务需求调整代码,增加更多功能如权限控制、统计分析等。

抱歉,评论功能暂时关闭!