如何用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
相关产品推荐
相关产品推荐

