StackSaga SQL 数据库分区支持 (SQL Database Partition Support)

概述 (Overview)

在使用 Saga 设计模式的高并发分布式系统中,事件存储表(主要是 es_transaction 和 es_transaction_execution_tryout)承受着极高频度的连续写入操作,因为系统需要详尽记录每一步事务状态与试算执行日志。

为了长期维持最优的读写查询性能、快速的索引构建以及可预测的磁盘空间开销,这些表均采用了范围分区 (Range Partitioning,例如基于 created_at 的按日自然区间分区)。

stacksaga-sql-partition-support 是一个轻量级、非阻塞的 Spring Boot Starter,旨在自动化创建和维护您整个生态中所有基于 SQL 的事件存储数据库的未来分区。

架构设计:集中式分区运行器 (Centralized Partition Runner)

在微服务架构中,每个业务服务都会引入各自对应的数据库响应式 Starter(stacksaga-mysql-reactive-support 或 stacksaga-pg-reactive-support)来执行分布式 Saga 事务。

全系统仅需部署一个中心实例:
与嵌入在每个微服务实例内部的运行时数据库支持模块不同,stacksaga-sql-partition-support 在整个系统中应当仅作为单个独立实例部署(或作为一个定时调度的独立运维任务)。

为什么推荐采用单一的集中式运行器架构?

  1. 避免多实例并发 DDL 冲突与锁争用: 如果每个微服务的每个 Pod 或副本在启动时或通过定时任务都尝试执行 ALTER TABLE …​ REORGANIZE PARTITION 等 Schema DDL,数据库将遭遇严重的元数据锁争用、连接数激增甚至是死锁。

  2. 严格的运行时数据访问权限边界: 应用微服务应以最小数据库权限运行(仅需 DML 权限:SELECT、INSERT、UPDATE、DELETE),绝不应赋予 DDL 表结构变更权限。通过集中式运行器,只有这单个运行器需要被授予必要的 ALTER / CREATE 权限。

  3. 跨多数据库与多引擎厂商的大一统归集: 即使您架构中的不同微服务使用了不同的数据库类型(例如订单服务使用 MySQL、支付服务使用 PostgreSQL、库存服务使用另一个独立的 MySQL),您也不需要为每个服务分别编写分区脚本。单个 stacksaga-sql-partition-support 实例可以配置多个数据源 (Multiple Datasources),跨越多种数据库引擎并并发统一处理。

StackSaga SQL 分区架构图

单点故障 (SPOF) 考量:为什么这完全不是问题

架构师常有的顾虑是:“整个系统仅有一个分区运行器实例,是否会引入单点故障 (SPOF)?”

答案是完全不会。在实际生产实践中,该设计具备极高的韧性弹性,对生产业务负载毫无风险:

  • 可配置的未来预建缓冲窗口 (unit + ahead-count): 当运行器触发时,它绝不仅仅创建当天的分区。通过 ahead-count 和 unit(DAYS、WEEK、MONTH、YEAR),它会提前预先创建未来多个周期的时间分区。 例如,配置 unit: DAYS 且 ahead-count: 7,或 unit: MONTH 且 ahead-count: 2,即可确保数据库中始终已经提前存在数天或数月之久的合法分区。

  • 与运行时业务事务完全解耦: 运行时的微服务(order-service、payment-service 等)直接向数据库插入 Saga 事件。在事务执行期间,它们与分区运行器之间*没有任何*网络通信或直接依赖。

  • 从容容忍运行器停机维护: 即便该单一分区运行器实例长时间下线、遭遇基础设施崩溃,或经历数天的停机维护,您的微服务依然能够毫无阻碍地持续写入事件,因为未来的分区表早已预先就绪在数据库中。

  • 重启后自动化追赶补偿: 一旦分区运行器重启或再次触发,它会自动探测数据库,发现任何缺失的未来时间区间,并幂等地完成增补建表。

核心特性 (Key Features)

  • 启动即检与定时调度执行: 在应用启动时立即检查并预建分区(run-on-startup=true),并支持通过可配置的 cron 表达式周期性执行(cron,默认为 0 0 12 * * * —— 每天中午 12:00)。

  • 多数据源并发处理: 支持在 stacksaga.partitioning.sql.datasources.* 下配置任意数量的数据源。所有配置的数据源均借助 Project Reactor 非阻塞的 Flux.flatMap(concurrency = 4) 并行并发执行。

  • 按需创建的非连接池连接 (Unpooled On-Demand Connections): 因为分区维护仅在应用启动和每日定时任务触发时运行数秒,常驻空闲连接池完全没有必要。该模块按需创建直连的无连接池 R2DBC 连接,执行完分区 DDL 后立即安全关闭,零资源浪费。

  • 方言自动推断 (MySQL 与 PostgreSQL): 根据 R2DBC URL 和连接元数据自动识别底层数据库类型:

    • MySQL:使用 p_max 兜底边界动态重组范围分区:

ALTER TABLE es_transaction REORGANIZE PARTITION p_max INTO (
    PARTITION p2026_09_13 VALUES LESS THAN ('2026-09-14 00:00:00'),
    PARTITION p_max VALUES LESS THAN (MAXVALUE)
);
  • PostgreSQL:创建声明式子分区:

CREATE TABLE IF NOT EXISTS es_transaction_p2026_09_13 PARTITION OF es_transaction
    FOR VALUES FROM ('2026-09-13 00:00:00') TO ('2026-09-14 00:00:00');
  • 启动前预检探活与权限校验: 在启动时验证数据库连通性(validate-on-startup=true),并通过 SHOW GRANTS (MySQL) 或 has_schema_privilege (PostgreSQL) 检查用户操作权限(validate-privileges=true),在缺少必要权限时尽早告警。

  • 提前预建机制 (ahead-count): 可针对每个数据源独立配置。设置 ahead-count: 2 会提前为今天以及未来 2 天预先创建分区,彻底杜绝跨过午夜零点时的数据写入拒绝异常。

引入依赖 (Adding as a Dependency)

在您的独立分区运行器微服务中引入 stacksaga-sql-partition-support:

<dependencyManagement>
    <dependencies>
        <dependency> <!--仅用于 stacksaga 依赖项版本管理-->
            <groupId>org.stacksaga</groupId>
            <artifactId>stacksaga-bom</artifactId>
            <version>1.0.0-SNAPSHOT</version>
            <type>pom</type>
            <scope>import</scope>
        </dependency>
    </dependencies>
</dependencyManagement>

<dependencies>
    <dependency>
        <groupId>org.stacksaga</groupId>
        <artifactId>stacksaga-sql-partition-support</artifactId>
    </dependency>
    <!-- 引入与所配置数据库匹配的响应式 R2DBC 驱动 -->
    <dependency>
        <groupId>io.asyncer</groupId>
        <artifactId>r2dbc-mysql</artifactId>
    </dependency>
    <dependency>
        <groupId>org.postgresql</groupId>
        <artifactId>r2dbc-postgresql</artifactId>
    </dependency>
</dependencies>

配置参数参考 (Configuration Reference)

Application Properties 配置示例

# 启用/禁用 starter (默认: true)
stacksaga.partitioning.sql.enabled=true

# 启动时校验数据库连通性 (默认: true)
stacksaga.partitioning.sql.validate-on-startup=true

# 启动时校验用户分区变更权限 (默认: true)
stacksaga.partitioning.sql.validate-privileges=true

# 启动时立即执行一次分区创建 (默认: true)
stacksaga.partitioning.sql.run-on-startup=true

# 定时调度 cron (默认: 每天中午 12:00)
stacksaga.partitioning.sql.cron=0 0 12 * * *

# 跨数据源的并发度限制 (默认: 4)
stacksaga.partitioning.sql.concurrency=4

# 建立 R2DBC 连接的超时时间 (默认: 10s)
stacksaga.partitioning.sql.connect-timeout=10s

# --- 数据源 1: 订单服务数据库 (MySQL) - 按天分区 ---
stacksaga.partitioning.sql.datasources.orderdb.url=r2dbc:mysql://localhost:3306/order_db
stacksaga.partitioning.sql.datasources.orderdb.username=partition_admin
stacksaga.partitioning.sql.datasources.orderdb.password=secret
stacksaga.partitioning.sql.datasources.orderdb.unit=DAYS
stacksaga.partitioning.sql.datasources.orderdb.ahead-count=2

# --- 数据源 2: 支付服务数据库 (PostgreSQL) - 按周分区 ---
stacksaga.partitioning.sql.datasources.paymentdb.url=r2dbc:postgresql://localhost:5432/payment_db
stacksaga.partitioning.sql.datasources.paymentdb.username=partition_admin
stacksaga.partitioning.sql.datasources.paymentdb.password=secret
stacksaga.partitioning.sql.datasources.paymentdb.unit=WEEK
stacksaga.partitioning.sql.datasources.paymentdb.ahead-count=3

# --- 数据源 3: 库存服务数据库 (MySQL) - 按月分区 ---
stacksaga.partitioning.sql.datasources.inventorydb.url=r2dbc:mysql://localhost:3306/inventory_db
stacksaga.partitioning.sql.datasources.inventorydb.username=partition_admin
stacksaga.partitioning.sql.datasources.inventorydb.password=secret
stacksaga.partitioning.sql.datasources.inventorydb.unit=MONTH
stacksaga.partitioning.sql.datasources.inventorydb.ahead-count=1

配置属性参考对照表

属性名称 (Property Name) 默认值 (Default) 类型 (Type) 描述说明 (Description)

stacksaga.partitioning.sql.enabled

true

boolean

启用或禁用分区支持 starter。

stacksaga.partitioning.sql.validate-on-startup

true

boolean

是否在应用启动时测试数据库连通性。

stacksaga.partitioning.sql.validate-privileges

true

boolean

是否在启动时验证数据库账号是否拥有 ALTER / CREATE 权限。

stacksaga.partitioning.sql.run-on-startup

true

boolean

是否在应用启动时立即执行一次分区创建。

stacksaga.partitioning.sql.exit-on-completion

false

boolean

启动分区创建完成后是否自动终止退出 JVM 应用程序(专为 Kubernetes CronJobs 设计)。所有数据源均成功时以退出码 0 退出,任何一个失败则以 1 退出。

stacksaga.partitioning.sql.cron

0 0 12 * * *

String

定时调度的 Spring cron 表达式(默认为每天中午 12:00)。当由外部调度器(如 Kubernetes CronJob)托管时,可设置为 none 或 - 以禁用内部定时器。

stacksaga.partitioning.sql.concurrency

4

int

并行并发处理的最大数据源数量。

stacksaga.partitioning.sql.connect-timeout

10s

Duration

建立按需 R2DBC 连接的超时时间。

stacksaga.partitioning.sql.datasources.<name>.url

-

String

R2DBC 连接 URL(例如 r2dbc:mysql://host:3306/db 或 r2dbc:postgresql://host:5432/db)。

stacksaga.partitioning.sql.datasources.<name>.username

-

String

具备分区 DDL 变更权限的数据库用户名。

stacksaga.partitioning.sql.datasources.<name>.password

-

String

数据库密码。

stacksaga.partitioning.sql.datasources.<name>.unit

DAYS

PartitionUnit

分区时间区间单位:DAYS(日)、WEEK(周)、MONTH(月)或 YEAR(年)。

stacksaga.partitioning.sql.datasources.<name>.ahead-count

1

int

提前预创建的分区区间数量。

stacksaga.partitioning.sql.datasources.<name>.type

-

DatabaseType

显式指定数据库类型(MYSQL 或 POSTGRESQL)。若省略,系统会根据 URL 和元数据自动识别。

部署策略 (Deployment Strategies)

因为分区维护仅需数秒钟即可完成所有校验与建表,您可以根据运维平台在以下两种部署策略中进行选择:

策略 A:持久化后台常驻守护进程 (传统虚拟机 / 容器)

在传统的虚拟机或标准 Docker 容器环境中:

  • 应用程序作为后台守护进程 7x24 小时常驻运行。

  • 启动时执行连通性探活、验证用户权限(validate-on-startup=true),并预先创建分区(run-on-startup=true)。

  • 启动执行完成后,在后台休眠,直到达到 stacksaga.partitioning.sql.cron(如 0 0 12 * * *)所指定的下一次定时调度时刻。

策略 B:定时瞬态作业 (Kubernetes CronJob)

在云原生容器编排环境(如 Kubernetes)中,保持一个 Pod 7x24 小时常驻运行并不经济,因为应用在数秒内建完分区后,整整一天都处于无所事事的空闲状态。

为了优化集群计算资源利用率,您可以将分区运行器部署为 Kubernetes CronJob(或 AWS ECS / GCP Cloud Run 上的 Scheduled Jobs):

  • 运行机制:

    1. Kubernetes 按照预定计划(例如每天中午 12:00)拉起 CronJob Pod。

    2. Pod 启动并初始化 Spring Boot,在 run-on-startup=true 的作用下,立即并发校验并为所有数据库创建未来的分区。

    3. 将 stacksaga.partitioning.sql.cron=none 设置为禁用内部 Spring 定时器。

    4. 开启 stacksaga.partitioning.sql.exit-on-completion=true,命令运行器在分区创建完成后自动退出 Spring Boot 与 JVM:

  • 若所有数据源均成功 → 干净退出,退出码为 0 (Completed)。

  • 若有任何数据源失败 → 以退出码 1 退出 (Error/Failed),触发 Kubernetes 按照您的 restartPolicy 进行告警或重试。

    1. Kubernetes 立即回收释放 Pod 的内存和 CPU 资源,直到下一个计划周期。

Kubernetes CronJob 清单示例:

apiVersion: batch/v1
kind: CronJob
metadata:
  name: stacksaga-partition-runner
  namespace: stacksaga
spec:
  schedule: "0 12 * * *"  # 每天中午 12:00 准时触发
  concurrencyPolicy: Forbid
  successfulJobsHistoryLimit: 3
  failedJobsHistoryLimit: 3
  jobTemplate:
    spec:
      template:
        spec:
          restartPolicy: OnFailure
          containers:
            - name: partition-runner
              image: your-docker-registry.internal/stacksaga/stacksaga-partition-runner:1.0.0
              env:
                # Pod 启动后立即执行分区创建
                - name: STACKSAGA_PARTITIONING_SQL_RUN_ON_STARTUP
                  value: "true"
                # 启动时校验连通性与权限
                - name: STACKSAGA_PARTITIONING_SQL_VALIDATE_ON_STARTUP
                  value: "true"
                - name: STACKSAGA_PARTITIONING_SQL_VALIDATE_PRIVILEGES
                  value: "true"
                # 禁用内部 cron,调度权交给 Kubernetes
                - name: STACKSAGA_PARTITIONING_SQL_CRON
                  value: "none"
                # 执行完毕后自动以退出码 0 (或错误时为 1) 退出
                - name: STACKSAGA_PARTITIONING_SQL_EXIT_ON_COMPLETION
                  value: "true"
              resources:
                requests:
                  memory: "256Mi"
                  cpu: "100m"
                limits:
                  memory: "512Mi"
                  cpu: "500m"

监控与可观测性 (Monitoring & Observability)

stacksaga-sql-partition-support 提供了可观测性生命周期钩子,用于监控分区执行结果、上报指标并在失败时触发报警。

PartitionExecutionCompletedEvent 事件

每当分区创建执行完毕时(无论是在应用启动时还是在定时触发时),运行器都会自动发布一个 Spring PartitionExecutionCompletedEvent 事件。

该事件提供了丰富的执行遥测数据:

方法名 返回值类型 描述说明

isAllSuccessful()

boolean

如果所有配置的数据源都成功完成分区且无任何异常,则返回 true。

getTotalDatasources()

int

本次执行所评估的数据源总数。

getSuccessCount()

int

成功完成分区的数据源数量。

getFailureCount()

int

分区失败的数据源数量。

getDuration()

Duration

执行消耗的总耗时。

getReports()

Map<String, DatasourceReport>

按数据源划分的详细执行报告,包含表名、分区数以及异常错误原因(若有)。

示例:监听分区事件并发送告警

开发人员可以注册一个 Spring @EventListener 来捕获分区指标,将失败告警推送到企业微信、钉钉、Slack 或 PagerDuty,或导出指标到 Prometheus:

import lombok.extern.slf4j.Slf4j;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Component;
import org.stacksaga.partition.event.PartitionExecutionCompletedEvent;

@Component
@Slf4j
public class PartitionMonitoringListener {

    @EventListener
    public void onPartitionCompleted(PartitionExecutionCompletedEvent event) {
        if (!event.isAllSuccessful()) {
            log.error("ALERT: 分区创建失败!共评估 {} 个数据源,失败 {} 个,耗时 {} ms!",
                    event.getTotalDatasources(), event.getFailureCount(), event.getDuration().toMillis());

            event.getReports().forEach((dsName, report) -> {
                if (!report.isSuccess()) {
                    log.error("数据源 '{}' 分区失败: {}", dsName, report.getErrorMessage(), report.getError());
                    // 示例:向钉钉、企微、Slack 发送告警通知
                    // alertService.sendAlert("数据库分区失败,数据源: " + dsName);
                }
            });
        } else {
            log.info("SUCCESS: 所有 {} 个数据源分区全部成功完成,总耗时 {} ms。",
                    event.getTotalDatasources(), event.getDuration().toMillis());
        }
    }
}