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

Flyway 5.0.7:如何阻止回调在数据库版本高于类路径最新迁移时执行并抛异常?

在Flyway 5.0.7版本里要实现这个需求,其实可以通过自定义回调或者扩展核心迁移类来搞定,下面给你两个具体的实现方案:

方案1:自定义AfterMigrate回调,检测版本后抛出异常

我们可以创建一个自定义的Flyway回调类,在afterMigrate事件触发前先检查数据库版本和本地最新迁移文件的版本差异,如果数据库版本更高,直接抛出异常,跳过原本的回调逻辑。

首先编写自定义回调类:

import org.flywaydb.core.api.callback.Callback;
import org.flywaydb.core.api.callback.Context;
import org.flywaydb.core.api.callback.Event;
import org.flywaydb.core.api.configuration.Configuration;
import org.flywaydb.core.internal.info.MigrationInfoServiceImpl;
import org.flywaydb.core.internal.metadatatable.MetaDataTable;
import org.flywaydb.core.internal.metadatatable.MetaDataTableRow;

import java.util.List;
import java.util.Comparator;

public class VersionCheckAfterMigrateCallback implements Callback {

    @Override
    public boolean supports(Event event, Context context) {
        // 只针对afterMigrate事件生效
        return Event.AFTER_MIGRATE.equals(event);
    }

    @Override
    public boolean canHandleInTransaction(Event event, Context context) {
        return false;
    }

    @Override
    public void handle(Event event, Context context) {
        Configuration config = context.getConfiguration();
        MetaDataTable metaDataTable = new MetaDataTable(
                config.getDataSource(), 
                config.getMetadataTableName(), 
                config.getSqlScriptFactory()
        );

        // 获取数据库中已应用的最新迁移版本
        List<MetaDataTableRow> appliedMigrations = metaDataTable.list();
        MetaDataTableRow latestDbMigration = appliedMigrations.stream()
                .max(Comparator.comparing(MetaDataTableRow::getVersion))
                .orElse(null);

        // 获取本地类路径里的最新迁移版本
        MigrationInfoServiceImpl migrationInfoService = new MigrationInfoServiceImpl(
                config, 
                config.getLocations(), 
                config.getResolvers()
        );
        migrationInfoService.refresh();
        String latestLocalVersion = migrationInfoService.getLatestAvailableMigration() != null
                ? migrationInfoService.getLatestAvailableMigration().getVersion().toString()
                : null;

        // 对比版本,数据库版本更高则抛出异常
        if (latestDbMigration != null && latestLocalVersion != null) {
            if (isDbVersionNewer(latestDbMigration.getVersion(), latestLocalVersion)) {
                throw new IllegalStateException(String.format(
                        "Schema \"%s\" has a version (%s) that is newer than the latest available migration (%s)!",
                        config.getSchemaName(), latestDbMigration.getVersion(), latestLocalVersion
                ));
            }
        }

        // 如果版本正常,这里可以添加你原本要执行的afterMigrate逻辑
        // ...
    }

    // 自定义版本比较逻辑,适配x.y.z格式的版本号
    private boolean isDbVersionNewer(String dbVersion, String localVersion) {
        String[] dbParts = dbVersion.split("\\.");
        String[] localParts = localVersion.split("\\.");

        int minLength = Math.min(dbParts.length, localParts.length);
        for (int i = 0; i < minLength; i++) {
            int dbNum = Integer.parseInt(dbParts[i]);
            int localNum = Integer.parseInt(localParts[i]);
            if (dbNum > localNum) {
                return true;
            } else if (dbNum < localNum) {
                return false;
            }
        }
        // 前面部分一致时,版本段更长的视为更高版本
        return dbParts.length > localParts.length;
    }
}

然后在配置Flyway时注册这个回调:

Flyway flyway = Flyway.configure()
        .dataSource(yourDataSource)
        .locations("db/migration") // 你的迁移文件路径
        .callbacks(new VersionCheckAfterMigrateCallback())
        .load();

flyway.migrate();
方案2:扩展DbMigrate核心类,提前拦截版本差异

Flyway的DbMigrate是执行迁移的核心类,它原本会在检测到数据库版本更高时打印警告。我们可以扩展这个类,在迁移执行前就做版本检查,直接抛出异常终止流程。

编写自定义的DbMigrate类:

import org.flywaydb.core.internal.command.DbMigrate;
import org.flywaydb.core.api.configuration.Configuration;
import org.flywaydb.core.internal.info.MigrationInfoServiceImpl;
import org.flywaydb.core.internal.metadatatable.MetaDataTable;
import java.util.Comparator;

public class VersionCheckingDbMigrate extends DbMigrate {

    public VersionCheckingDbMigrate(Configuration configuration) {
        super(configuration);
    }

    @Override
    protected void migrate() {
        Configuration config = getConfiguration();
        MetaDataTable metaDataTable = new MetaDataTable(
                config.getDataSource(),
                config.getMetadataTableName(),
                config.getSqlScriptFactory()
        );

        // 获取数据库最新版本
        String latestDbVersion = metaDataTable.getAppliedMigrations().stream()
                .max(Comparator.comparing(row -> row.getVersion()))
                .map(row -> row.getVersion())
                .orElse(null);

        // 获取本地最新迁移版本
        MigrationInfoServiceImpl migrationInfoService = new MigrationInfoServiceImpl(
                config,
                config.getLocations(),
                config.getResolvers()
        );
        migrationInfoService.refresh();
        String latestLocalVersion = migrationInfoService.getLatestAvailableMigration() != null
                ? migrationInfoService.getLatestAvailableMigration().getVersion().toString()
                : null;

        // 版本检查,不符合则抛出异常
        if (latestDbVersion != null && latestLocalVersion != null) {
            if (isDbVersionNewer(latestDbVersion, latestLocalVersion)) {
                throw new IllegalStateException(String.format(
                        "Schema \"%s\" has a version (%s) that is newer than the latest available migration (%s)!",
                        config.getSchemaName(), latestDbVersion, latestLocalVersion
                ));
            }
        }

        // 版本正常则执行原迁移逻辑
        super.migrate();
    }

    // 复用版本比较方法
    private boolean isDbVersionNewer(String dbVersion, String localVersion) {
        String[] dbParts = dbVersion.split("\\.");
        String[] localParts = localVersion.split("\\.");

        int minLength = Math.min(dbParts.length, localParts.length);
        for (int i = 0; i < minLength; i++) {
            int dbNum = Integer.parseInt(dbParts[i]);
            int localNum = Integer.parseInt(localParts[i]);
            if (dbNum > localNum) {
                return true;
            } else if (dbNum < localNum) {
                return false;
            }
        }
        return dbParts.length > localParts.length;
    }
}

使用时替换默认的DbMigrate执行迁移:

Configuration flywayConfig = Flyway.configure()
        .dataSource(yourDataSource)
        .locations("db/migration")
        .configuration();

VersionCheckingDbMigrate customMigrate = new VersionCheckingDbMigrate(flywayConfig);
customMigrate.migrate();
注意事项
  • 版本比较逻辑是基于标准的x.y.z格式编写的,如果你的版本包含预发布标签(比如11.2.5-SNAPSHOT),需要额外调整比较逻辑。
  • Flyway 5.0.7的内部API和后续版本有差异,以上代码仅适配该版本,升级Flyway后可能需要修改实现。
  • 如果使用多个Flyway回调,要注意回调的执行顺序,确保版本检查逻辑优先执行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:46:36