You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PostgreSQL自动更新Section状态并向Spring应用推送事件方案咨询

投票板块状态自动更新与事件通知解决方案

问题背景

开发投票板块应用时,需实现Section表的status字段根据时间自动变更:

  • WAITING_TO_START:当前时间早于start_date
  • IN_PROGRESS:当前时间在start_date与start_date+duration之间
  • FINISHED:当前时间晚于start_date+duration
    需求是在时间触发状态变更时,自动更新PostgreSQL中的status字段,并向Spring应用发送事件通知客户端。

现有代码与结构

Section表结构

Section表结构

Section Model代码

package com.testdbserver.desafiovotacao.data.models;

import com.fasterxml.jackson.annotation.JsonBackReference;
import com.testdbserver.desafiovotacao.data.enums.SectionStatusEnum;
import jakarta.persistence.*;
import jakarta.validation.constraints.FutureOrPresent;
import lombok.*;
import java.util.Date;
import java.util.UUID;

@NoArgsConstructor
@AllArgsConstructor
@Getter
@Setter
@Builder
@Entity
@Table(name="section")
public class Section {
    @Id
    @GeneratedValue(strategy = GenerationType.UUID)
    private UUID id;

    @ManyToOne
    @JoinColumn(name="pauta_id")
    private Pauta pauta;

    @Column(name="status", nullable = false)
    @Enumerated(EnumType.STRING)
    private SectionStatusEnum status;

    @Column(name="created_at", nullable = false)
    @Temporal(TemporalType.TIMESTAMP)
    private Date createdAt;

    @FutureOrPresent
    @Column(name="dt_start", nullable = false)
    @Temporal(TemporalType.TIMESTAMP)
    private Date dtStart;

    @Column(name="duration", nullable = false)
    private int duration;

    @JsonBackReference
    public Pauta getPauta() {
        return pauta;
    }
}

Section Repository代码

package com.testdbserver.desafiovotacao.data.repositories;

import com.testdbserver.desafiovotacao.data.models.Section;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.data.jpa.repository.Query;

import java.util.Date;
import java.util.List;
import java.util.UUID;

public interface SectionRepository extends JpaRepository<Section, UUID> {
    @Query(value = " SELECT s FROM Section s " +
            "         WHERE (:allowFinishedSections=false and s.status <> 'FINISHED') " +
            "           AND s.dt_start < :maxDate " +
            "         ORDER BY s.dt_start asc", nativeQuery = true)
    public List<Section> getAllSections(boolean allowFinishedSections, Date maxDate);
}

POM.XML配置

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
    xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <parent>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-parent</artifactId>
        <version>3.0.4</version>
        <relativePath/> <!-- lookup parent from repository -->
    </parent>
    <groupId>com.testdbserver</groupId>
    <artifactId>desafio-votacao</artifactId>
    <version>0.0.1-SNAPSHOT</version>
    <name>desafio-votacao</name>
    <description>Challenge project to create a vote system to a credit cooperative</description>
    <properties>
        <java.version>17</java.version>
    </properties>
    <dependencies>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-data-jpa</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-web</artifactId>
        </dependency>

        <dependency>
            <groupId>org.postgresql</groupId>
            <artifactId>postgresql</artifactId>
            <scope>runtime</scope>
        </dependency>
        <dependency>
            <groupId>org.projectlombok</groupId>
            <artifactId>lombok</artifactId>
            <optional>true</optional>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-test</artifactId>
            <scope>test</scope>
        </dependency>
        <dependency>
            <groupId>org.springdoc</groupId>
            <artifactId>springdoc-openapi-ui</artifactId>
            <version>1.6.15</version>
        </dependency>
        <dependency>
            <groupId>io.jsonwebtoken</groupId>
            <artifactId>jjwt-api</artifactId>
            <version>0.11.5</version>
        </dependency>
        <dependency>
            <groupId>io.jsonwebtoken</groupId>
            <artifactId>jjwt-impl</artifactId>
            <version>0.11.5</version>
        </dependency>
        <dependency>
            <groupId>io.jsonwebtoken</groupId>
            <artifactId>jjwt-jackson</artifactId>
            <version>0.11.5</version>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-security</artifactId>
        </dependency>
        <dependency>
            <groupId>org.hibernate</groupId>
            <artifactId>hibernate-validator</artifactId>
            <version>8.0.0.Final</version>
        </dependency>
    </dependencies>

    <build>
        <plugins>
            <plugin>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-maven-plugin</artifactId>
                <configuration>
                    <excludes>
                        <exclude>
                            <groupId>org.projectlombok</groupId>
                            <artifactId>lombok</artifactId>
                        </exclude>
                    </excludes>
                </configuration>
            </plugin>
        </plugins>
    </build>

</project>

解决方案

方案一:数据库层面自动更新+触发应用事件

1. PostgreSQL定时任务与触发器

  • 先安装pg_cron扩展(PostgreSQL 12+支持,需在postgresql.conf中启用)
  • 创建定时任务,定期检查并更新状态:
-- 每分钟执行一次(可按需调整频率)
SELECT cron.schedule('update-section-status', '* * * * *', $$
    UPDATE section
    SET status = CASE
        WHEN NOW() < dt_start THEN 'WAITING_TO_START'
        WHEN NOW() BETWEEN dt_start AND dt_start + INTERVAL '1 minute' * duration THEN 'IN_PROGRESS'
        ELSE 'FINISHED'
    END
    WHERE status != CASE
        WHEN NOW() < dt_start THEN 'WAITING_TO_START'
        WHEN NOW() BETWEEN dt_start AND dt_start + INTERVAL '1 minute' * duration THEN 'IN_PROGRESS'
        ELSE 'FINISHED'
    END;
$$);
  • 创建触发器,状态更新时向应用发送通知:
-- 定义通知函数
CREATE OR REPLACE FUNCTION notify_section_status_change()
RETURNS TRIGGER AS $$
BEGIN
    PERFORM pg_notify('section_status_update', row_to_json(NEW)::text);
    RETURN NEW;
END;
$$ LANGUAGE plpgsql;

-- 绑定触发器到section表的UPDATE事件
CREATE TRIGGER section_status_trigger
AFTER UPDATE OF status ON section
FOR EACH ROW
EXECUTE FUNCTION notify_section_status_change();

2. Spring应用监听数据库通知

@Component
public class SectionStatusChangeListener {

    @Autowired
    private EntityManager entityManager;

    @Autowired
    private ApplicationEventPublisher eventPublisher;

    @PostConstruct
    public void listenToNotifications() {
        Session session = entityManager.unwrap(Session.class);
        session.doWork(connection -> {
            try (Statement statement = connection.createStatement()) {
                statement.execute("LISTEN section_status_update");
            }
            // 循环监听消息
            while (true) {
                connection.getNotifications(1000);
                List<PGNotification> notifications = ((PGConnection) connection).getNotifications();
                if (notifications != null) {
                    for (PGNotification notification : notifications) {
                        // 解析消息并发布Spring事件
                        Section updatedSection = new ObjectMapper().readValue(notification.getParameter(), Section.class);
                        eventPublisher.publishEvent(new SectionStatusChangedEvent(updatedSection));
                    }
                }
            }
        });
    }
}

// 自定义状态变更事件类
public class SectionStatusChangedEvent extends ApplicationEvent {
    public SectionStatusChangedEvent(Section section) {
        super(section);
    }

    public Section getSection() {
        return (Section) getSource();
    }
}

方案二:Spring应用层面定时任务+事件发布

1. 定时任务检查并更新状态

@Component
@EnableScheduling
public class SectionStatusScheduler {

    @Autowired
    private SectionRepository sectionRepository;

    @Autowired
    private ApplicationEventPublisher eventPublisher;

    // 每分钟执行一次
    @Scheduled(fixedRate = 60000)
    public void updateSectionStatuses() {
        Date now = new Date();
        List<Section> sections = sectionRepository.findAll();
        for (Section section : sections) {
            SectionStatusEnum newStatus = determineStatus(section, now);
            if (!newStatus.equals(section.getStatus())) {
                section.setStatus(newStatus);
                Section updatedSection = sectionRepository.save(section);
                // 发布状态变更事件
                eventPublisher.publishEvent(new SectionStatusChangedEvent(updatedSection));
            }
        }
    }

    private SectionStatusEnum determineStatus(Section section, Date now) {
        Date startDate = section.getDtStart();
        Date endDate = new Date(startDate.getTime() + (long) section.getDuration() * 60 * 1000);
        if (now.before(startDate)) {
            return SectionStatusEnum.WAITING_TO_START;
        } else if (now.after(endDate)) {
            return SectionStatusEnum.FINISHED;
        } else {
            return SectionStatusEnum.IN_PROGRESS;
        }
    }
}

2. 事件通知客户端(WebSocket示例)

@Configuration
@EnableWebSocketMessageBroker
public class WebSocketConfig implements WebSocketMessageBrokerConfigurer {
    @Override
    public void configureMessageBroker(MessageBrokerRegistry config) {
        config.enableSimpleBroker("/topic");
        config.setApplicationDestinationPrefixes("/app");
    }

    @Override
    public void registerStompEndpoints(StompEndpointRegistry registry) {
        registry.addEndpoint("/section-status").withSockJS();
    }
}

// 事件监听器,推送状态更新到客户端
@Component
public class SectionStatusEventListener {

    @Autowired
    private SimpMessagingTemplate messagingTemplate;

    @EventListener
    public void handleSectionStatusChange(SectionStatusChangedEvent event) {
        messagingTemplate.convertAndSend("/topic/section-status-updates", event.getSection());
    }
}

方案选择建议

  • 多实例部署场景优先选数据库层面方案,避免重复执行任务;
  • 逻辑集中管理优先选Spring定时任务方案,更易维护;
  • 客户端实时通知优先用WebSocket,适合双向通信场景。

内容的提问来源于stack exchange,提问作者Eduardo Lopes

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.26 20:17:09