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

如何用Apache Camel更新JPA/Hibernate实体?批量更新异常求助

Apache Camel JPA路由最后一批记录未更新问题

项目环境

基于Spring Boot的Apache Camel项目,使用Maven依赖:camel-spring-boot-starter、camel-jpa-starter、camel-endpointdsl

JPA实体定义

项目包含三个关联实体:

@Entity
@Table(name = RawDataDelivery.TABLE_NAME)
@BatchSize(size = 10)
public class RawDataDelivery extends PersistentObjectWithCreationDate {
    protected static final String TABLE_NAME = "raw_data_delivery";

    private static final String COLUMN_CONFIGURATION_ID = "configuration_id";
    private static final String COLUMN_SCOPED_CALCULATED = "scopes_calculated";

    @Column(nullable = false, name = COLUMN_SCOPED_CALCULATED)
    private boolean scopesCalculated;

    @OneToMany(mappedBy = "raw_data_delivery", fetch = FetchType.LAZY)
    private Set<RawDataFile> files = new HashSet<>();

    @CollectionTable(name = "processed_scopes_per_delivery")
    @ElementCollection(targetClass = String.class)
    private Set<String> processedScopes = new HashSet<>();
    
    // Getter/Setter
}

@Entity
@Table(name = RawDataFile.TABLE_NAME)
@BatchSize(size = 100)
public class RawDataFile extends PersistentObjectWithCreationDate {
    protected static final String TABLE_NAME = "raw_data_files";

    private static final String COLUMN_CONFIGURATION_ID = "configuration_id";
    private static final String COLUMN_RAW_DATA_DELIVERY_ID = "raw_data_delivery_id";
    private static final String COLUMN_PARENT_ID = "parent_file_id";
    private static final String COLUMN_IDENTIFIER = "identifier";
    private static final String COLUMN_CONTENT = "content";
    private static final String COLUMN_FILE_SIZE_IN_BYTES = "file_size_in_bytes";

    @ManyToOne(optional = true, fetch = FetchType.LAZY)
    @JoinColumn(name = COLUMN_RAW_DATA_DELIVERY_ID)
    private RawDataDelivery rawDataDelivery;

    @Column(name = COLUMN_IDENTIFIER, nullable = false)
    private String identifier;

    @Lob
    @Column(name = COLUMN_CONTENT, nullable = true)
    private Blob content;

    @Column(name = COLUMN_FILE_SIZE_IN_BYTES, nullable = false)
    private long fileSizeInBytes;
    
    // Getter/Setter
}

@Entity
@TypeDef(name = "jsonb", typeClass = JsonBinaryType.class)
@Table(name = RawDataRecord.TABLE_NAME, uniqueConstraints = ...)
public class RawDataRecord extends PersistentObjectWithCreationDate {
    public static final String TABLE_NAME = "raw_data_records";

    static final String COLUMN_RAW_DATA_FILE_ID = "raw_data_file_id";
    static final String COLUMN_INDEX = "index";
    static final String COLUMN_CONTENT = "content";
    static final String COLUMN_HASHCODE = "hashcode";
    static final String COLUMN_SCOPE = "scope";

    @ManyToOne(optional = false)
    @JoinColumn(name = COLUMN_RAW_DATA_FILE_ID)
    private RawDataFile rawDataFile;

    @Column(name = COLUMN_INDEX, nullable = false)
    private long index;

    @Lob
    @Type(type = "jsonb")
    @Column(name = COLUMN_CONTENT, nullable = false, columnDefinition = "jsonb")
    private String content;

    @Column(name = COLUMN_HASHCODE, nullable = false)
    private String hashCode;

    @Column(name = COLUMN_SCOPE, nullable = true)
    private String scope;
}

需求与现有路由

需求是构建一条Camel路由:筛选scopesCalculated为false的RawDataDelivery记录,在同一PostgreSQL事务中更新其关联文件下所有RawDataRecord的scope字段,全部更新完成后将RawDataDelivery的scopesCalculated设为true并提交事务。

现有路由代码:

String r3RouteId = ...;

var dataSource3 = jpa(RawDataDelivery.class.getName())
        .lockModeType(LockModeType.NONE)
        .delay(60).timeUnit(TimeUnit.SECONDS)
        .consumeDelete(false)
        .query("select rdd from RawDataDelivery rdd where rdd.scopesCalculated is false and rdd.configuration.id = " + configuration.getId())
;

from(dataSource3)
        .routeId(r3RouteId)
        .routeDescription(configuration.getName())
        .messageHistory()
        .transacted()
        .process(exchange -> {
            RawDataDelivery rawDataDelivery = exchange.getIn().getBody(RawDataDelivery.class);
            rawDataDelivery.setScopesCalculated(true);
        })
        .transform(new Expression() {
            @Override
            public <T> T evaluate(Exchange exchange, Class<T> type) {
                RawDataDelivery rawDataDelivery = exchange.getIn().getBody(RawDataDelivery.class);
                return (T)rawDataDelivery.getFiles();
            }
        })
        .split(bodyAs(Iterator.class)).streaming()
            .transform(new Expression() {
                @Override
                public <T> T evaluate(Exchange exchange, Class<T> type) {
                    RawDataFile rawDataFile = exchange.getIn().getBody(RawDataFile.class);
                    
                    // rawDataRecordJpaRepository是Spring自动注入的接口,包含方法:
                    // @Lock(value = LockModeType.NONE)
                    // Stream<RawDataRecord> findByRawDataFile(RawDataFile rawDataFile);
                    
                    // 每个文件可能有10万+记录,为避免内存溢出,采用流式查询并按500条分批处理
                    return (T)rawDataRecordJpaRepository.findByRawDataFile(rawDataFile);
                }
            })
            .split(bodyAs(Iterator.class)).streaming()
                .aggregate(constant("all"), new GroupedBodyAggregationStrategy())
                    .completionSize(500)
                    .completionTimeout(TimeUnit.SECONDS.toMillis(5))
                    .process(exchange -> {
                        List<RawDataRecord> rawDataRecords = exchange.getIn().getBody(List.class);

                        for (RawDataRecord rawDataRecord : rawDataRecords) {
                            rawDataRecord.setScope("abc");
                        }
                    })
;

问题现象

路由大体可运行,但存在最后一批记录未更新的问题:比如总共有43782条记录,仅43500条被更新,剩余282条的scope仍为null。同时对Camel的JPA事务与会话管理机制不熟悉,找不到不用SQL组件的实体更新示例。

已尝试的无效方案

  • 设置jpa(...).joinTransaction(false).advanced().sharedEntityManager(true)
  • 使用.enrich(jpa(RawDataRecord.class.getName()).query("select rec from RawDataRecord rec where rawDataFile = ${body}"))替代JPA Repository的transform操作
  • 从Camel消息头获取Hibernate Session更新实体:Session session = exchange.getIn().getHeader(JpaConstants.ENTITY_MANAGER, Session.class);
  • 在路由末尾通过新的JPA组件执行更新:
.split(bodyAs(Iterator.class)).streaming()
.to(jpa(RawDataRecord.class.getName()).usePersist(false).flushOnSend(false))

求助

请问还有其他解决思路或建议吗?

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 07:15:38