From b76dffb9817504612190a1f27dd5d4a54f18ad8f Mon Sep 17 00:00:00 2001 From: Enrico Date: Wed, 6 May 2026 14:33:29 +0200 Subject: [PATCH] feat: decouple ChangeLogAccumulator from Hibernate Session [DHIS2-21378] --- .../persister/AbstractTrackerPersister.java | 196 ++++++++++-------- .../persister/ChangeLogAccumulator.java | 8 +- .../bundle/persister/EnrollmentPersister.java | 4 +- .../persister/RelationshipPersister.java | 5 +- .../persister/SingleEventPersister.java | 5 +- .../persister/TrackedEntityPersister.java | 5 +- .../persister/TrackerEventPersister.java | 5 +- 7 files changed, 126 insertions(+), 102 deletions(-) diff --git a/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/AbstractTrackerPersister.java b/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/AbstractTrackerPersister.java index 8563241854a1..25712dbdbc42 100644 --- a/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/AbstractTrackerPersister.java +++ b/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/AbstractTrackerPersister.java @@ -36,6 +36,8 @@ import jakarta.persistence.EntityManager; import jakarta.persistence.PersistenceContext; +import java.sql.Connection; +import java.sql.SQLException; import java.time.ZoneId; import java.time.format.DateTimeFormatter; import java.util.ArrayList; @@ -49,6 +51,7 @@ import java.util.Set; import java.util.function.Function; import java.util.stream.Collectors; +import javax.sql.DataSource; import lombok.AccessLevel; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; @@ -78,6 +81,7 @@ import org.hisp.dhis.tracker.model.TrackedEntity; import org.hisp.dhis.tracker.model.TrackedEntityAttributeValue; import org.hisp.dhis.user.UserDetails; +import org.springframework.jdbc.datasource.DataSourceUtils; /** * @author Luciano Fiandesio @@ -93,6 +97,8 @@ public abstract class AbstractTrackerPersister dtos = getByType(bundle); - for (T trackerDto : dtos) { - - Entity objectReport = new Entity(getType(), trackerDto.getUID()); - boolean isNewEntity = isNew(bundle, trackerDto); - // Capture before convert() which mutates the preheat entity's status - boolean completedInThisImport = - !bundle.isSkipSideEffects() - && isBeingCompleted(bundle.getPreheat(), trackerDto, isNewEntity); - ChangeLogAccumulator.Mark mark = changeLogs.mark(); - try { - V originalEntity = cloneEntityProperties(bundle.getPreheat(), trackerDto); - - // - // Convert the TrackerDto into an Hibernate-managed entity - // - V convertedDto = convert(bundle, trackerDto); - - // - // Handle ownership records, if required - // - persistOwnership(bundle, trackerDto, convertedDto); - - // - // Save or update the entity - // - if (isNew(bundle, trackerDto)) { - entityManager.persist(convertedDto); - updateDataValues( - bundle.getPreheat(), - trackerDto, - convertedDto, - originalEntity, - bundle.getUser(), - changeLogs); - typeReport.getStats().incCreated(); - typeReport.addEntity(objectReport); - updateAttributes( - bundle.getPreheat(), trackerDto, convertedDto, bundle.getUser(), changeLogs); - bundle.addUpdatedTrackedEntities(getUpdatedTrackedEntities(convertedDto)); - } else { - if (trackerDto.getTrackerType() == TrackerType.RELATIONSHIP) { - typeReport.getStats().incIgnored(); - // Relationships are not updated. A warning was already added to the report - } else { + Connection conn = DataSourceUtils.getConnection(dataSource); + try { + for (T trackerDto : dtos) { + + Entity objectReport = new Entity(getType(), trackerDto.getUID()); + boolean isNewEntity = isNew(bundle, trackerDto); + // Capture before convert() which mutates the preheat entity's status + boolean completedInThisImport = + !bundle.isSkipSideEffects() + && isBeingCompleted(bundle.getPreheat(), trackerDto, isNewEntity); + ChangeLogAccumulator.Mark mark = changeLogs.mark(); + try { + V originalEntity = cloneEntityProperties(bundle.getPreheat(), trackerDto); + + // + // Convert the TrackerDto into an Hibernate-managed entity + // + V convertedDto = convert(bundle, trackerDto); + + // + // Handle ownership records, if required + // + persistOwnership(bundle, trackerDto, convertedDto); + + // + // Save or update the entity + // + if (isNew(bundle, trackerDto)) { + entityManager.persist(convertedDto); updateDataValues( bundle.getPreheat(), trackerDto, @@ -166,62 +157,91 @@ public PersistResult persist(TrackerBundle bundle) { originalEntity, bundle.getUser(), changeLogs); + typeReport.getStats().incCreated(); + typeReport.addEntity(objectReport); updateAttributes( bundle.getPreheat(), trackerDto, convertedDto, bundle.getUser(), changeLogs); - entityManager.merge(convertedDto); - typeReport.getStats().incUpdated(); - typeReport.addEntity(objectReport); bundle.addUpdatedTrackedEntities(getUpdatedTrackedEntities(convertedDto)); + } else { + if (trackerDto.getTrackerType() == TrackerType.RELATIONSHIP) { + typeReport.getStats().incIgnored(); + // Relationships are not updated. A warning was already added to the report + } else { + updateDataValues( + bundle.getPreheat(), + trackerDto, + convertedDto, + originalEntity, + bundle.getUser(), + changeLogs); + updateAttributes( + bundle.getPreheat(), trackerDto, convertedDto, bundle.getUser(), changeLogs); + entityManager.merge(convertedDto); + typeReport.getStats().incUpdated(); + typeReport.addEntity(objectReport); + bundle.addUpdatedTrackedEntities(getUpdatedTrackedEntities(convertedDto)); + } } - } - if (!bundle.isSkipSideEffects()) { - EntityNotifications entityNotifications = - collectNotifications(bundle, convertedDto, isNewEntity, completedInThisImport); - if (entityNotifications != null) { - notifications.add(entityNotifications); + if (!bundle.isSkipSideEffects()) { + EntityNotifications entityNotifications = + collectNotifications(bundle, convertedDto, isNewEntity, completedInThisImport); + if (entityNotifications != null) { + notifications.add(entityNotifications); + } } - } - - // - // Add the entity to the Preheat - // - updatePreheat(bundle.getPreheat(), convertedDto); - if (FlushMode.OBJECT == bundle.getFlushMode()) { - // Flush entity INSERTs/UPDATEs before changelog INSERTs so FK references - // (trackedentityid, eventid) exist before changelog rows reference them. - entityManager.flush(); - changeLogs.flushAll(entityManager); - } - } catch (Exception e) { - changeLogs.rollbackTo(mark); - - final String msg = - "A Tracker Entity of type '" - + getType().getName() - + "' (" - + trackerDto.getUID() - + ") failed to persist."; - - if (AtomicMode.ALL.equals(bundle.getAtomicMode())) { - throw new PersistenceException(msg, e); - } else { - // TODO currently we do not keep track of the failed entity - // in the TrackerObjectReport + // + // Add the entity to the Preheat + // + updatePreheat(bundle.getPreheat(), convertedDto); - log.warn(msg + "\nThe Import process will process remaining entities.", e); + if (FlushMode.OBJECT == bundle.getFlushMode()) { + // Flush entity INSERTs/UPDATEs before changelog INSERTs so FK references + // (trackedentityid, eventid) exist before changelog rows reference them. + entityManager.flush(); + changeLogs.flushAll(conn); + } + } catch (Exception e) { + changeLogs.rollbackTo(mark); + + final String msg = + "A Tracker Entity of type '" + + getType().getName() + + "' (" + + trackerDto.getUID() + + ") failed to persist."; + + if (AtomicMode.ALL.equals(bundle.getAtomicMode())) { + throw new PersistenceException(msg, e); + } else { + // TODO currently we do not keep track of the failed entity + // in the TrackerObjectReport + + // TODO: if the failure originated from a JDBC flush (changeLogs.flushAll or, from + // Phase 3 onward, EntityWriteBatch.flush), the underlying PostgreSQL connection is + // now in an aborted-transaction state. Any subsequent SQL on the same connection + // will fail with "current transaction is aborted". This means a single JDBC flush + // failure in non-atomic mode silently cascades and causes all remaining entities to + // be ignored as well. Consider wrapping each entity's JDBC flush in a savepoint so + // that a failure can be rolled back to the savepoint and the connection stays usable. + log.warn(msg + "\nThe Import process will process remaining entities.", e); - typeReport.getStats().incIgnored(); + typeReport.getStats().incIgnored(); + } } } - } - if (FlushMode.AUTO == bundle.getFlushMode()) { - // Flush entity INSERTs/UPDATEs before changelog INSERTs so FK references - // (trackedentityid, eventid) exist before changelog rows reference them. - entityManager.flush(); - changeLogs.flushAll(entityManager); + if (FlushMode.AUTO == bundle.getFlushMode()) { + // Flush entity INSERTs/UPDATEs before changelog INSERTs so FK references + // (trackedentityid, eventid) exist before changelog rows reference them. + entityManager.flush(); + changeLogs.flushAll(conn); + } + } catch (SQLException e) { + throw new PersistenceException(e); + } finally { + DataSourceUtils.releaseConnection(conn, dataSource); } return new PersistResult(typeReport, notifications); } diff --git a/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/ChangeLogAccumulator.java b/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/ChangeLogAccumulator.java index 4ba399dfc391..4761952a9850 100644 --- a/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/ChangeLogAccumulator.java +++ b/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/ChangeLogAccumulator.java @@ -29,7 +29,6 @@ */ package org.hisp.dhis.tracker.imports.bundle.persister; -import jakarta.persistence.EntityManager; import java.sql.Connection; import java.sql.PreparedStatement; import java.sql.SQLException; @@ -40,7 +39,6 @@ import java.util.List; import javax.annotation.CheckForNull; import javax.annotation.Nonnull; -import org.hibernate.Session; import org.hisp.dhis.changelog.ChangeLogType; import org.hisp.dhis.dataelement.DataElement; import org.hisp.dhis.program.Program; @@ -185,15 +183,15 @@ void rollbackTo(Mark mark) { truncate(singleEventChangeLogs, mark.singleEventSize); } - void flushAll(EntityManager entityManager) { + void flushAll(Connection connection) throws SQLException { if (teChangeLogs.isEmpty() && trackerEventChangeLogs.isEmpty() && singleEventChangeLogs.isEmpty()) { return; } - Session session = entityManager.unwrap(Session.class); - session.doWork(this::insertAll); + insertAll(connection); + teChangeLogs.clear(); trackerEventChangeLogs.clear(); singleEventChangeLogs.clear(); diff --git a/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/EnrollmentPersister.java b/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/EnrollmentPersister.java index 1ba0901cec20..8f796a3e5f4c 100644 --- a/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/EnrollmentPersister.java +++ b/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/EnrollmentPersister.java @@ -32,6 +32,7 @@ import java.util.EnumSet; import java.util.List; import java.util.Set; +import javax.sql.DataSource; import org.hisp.dhis.common.UID; import org.hisp.dhis.program.EnrollmentStatus; import org.hisp.dhis.program.notification.NotificationTrigger; @@ -58,8 +59,9 @@ public class EnrollmentPersister public EnrollmentPersister( ReservedValueService reservedValueService, + DataSource dataSource, TrackedEntityProgramOwnerService trackedEntityProgramOwnerService) { - super(reservedValueService); + super(reservedValueService, dataSource); this.trackedEntityProgramOwnerService = trackedEntityProgramOwnerService; } diff --git a/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/RelationshipPersister.java b/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/RelationshipPersister.java index 309ac3628ff6..aa96d4bc7df2 100644 --- a/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/RelationshipPersister.java +++ b/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/RelationshipPersister.java @@ -31,6 +31,7 @@ import java.util.List; import java.util.Set; +import javax.sql.DataSource; import org.hisp.dhis.common.UID; import org.hisp.dhis.reservedvalue.ReservedValueService; import org.hisp.dhis.tracker.TrackerType; @@ -49,8 +50,8 @@ public class RelationshipPersister extends AbstractTrackerPersister { - public RelationshipPersister(ReservedValueService reservedValueService) { - super(reservedValueService); + public RelationshipPersister(ReservedValueService reservedValueService, DataSource dataSource) { + super(reservedValueService, dataSource); } @Override diff --git a/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/SingleEventPersister.java b/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/SingleEventPersister.java index f69c39068d1a..db1917d65a31 100644 --- a/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/SingleEventPersister.java +++ b/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/SingleEventPersister.java @@ -44,6 +44,7 @@ import java.util.stream.Collectors; import javax.annotation.CheckForNull; import javax.annotation.Nonnull; +import javax.sql.DataSource; import org.apache.commons.lang3.StringUtils; import org.hisp.dhis.changelog.ChangeLogType; import org.hisp.dhis.common.UID; @@ -73,8 +74,8 @@ public class SingleEventPersister extends AbstractTrackerPersister< org.hisp.dhis.tracker.imports.domain.SingleEvent, SingleEvent> { - public SingleEventPersister(ReservedValueService reservedValueService) { - super(reservedValueService); + public SingleEventPersister(ReservedValueService reservedValueService, DataSource dataSource) { + super(reservedValueService, dataSource); } @Override diff --git a/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/TrackedEntityPersister.java b/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/TrackedEntityPersister.java index 4917d8985b06..18babf87bf85 100644 --- a/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/TrackedEntityPersister.java +++ b/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/TrackedEntityPersister.java @@ -32,6 +32,7 @@ import java.util.Collections; import java.util.List; import java.util.Set; +import javax.sql.DataSource; import org.hisp.dhis.common.UID; import org.hisp.dhis.reservedvalue.ReservedValueService; import org.hisp.dhis.tracker.TrackerType; @@ -50,8 +51,8 @@ public class TrackedEntityPersister extends AbstractTrackerPersister< org.hisp.dhis.tracker.imports.domain.TrackedEntity, TrackedEntity> { - public TrackedEntityPersister(ReservedValueService reservedValueService) { - super(reservedValueService); + public TrackedEntityPersister(ReservedValueService reservedValueService, DataSource dataSource) { + super(reservedValueService, dataSource); } @Override diff --git a/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/TrackerEventPersister.java b/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/TrackerEventPersister.java index 31b6f3e66977..4394bc349a11 100644 --- a/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/TrackerEventPersister.java +++ b/dhis-2/dhis-tracker/src/main/java/org/hisp/dhis/tracker/imports/bundle/persister/TrackerEventPersister.java @@ -45,6 +45,7 @@ import java.util.stream.Stream; import javax.annotation.CheckForNull; import javax.annotation.Nonnull; +import javax.sql.DataSource; import org.apache.commons.lang3.StringUtils; import org.hisp.dhis.changelog.ChangeLogType; import org.hisp.dhis.common.UID; @@ -74,8 +75,8 @@ public class TrackerEventPersister extends AbstractTrackerPersister< org.hisp.dhis.tracker.imports.domain.TrackerEvent, TrackerEvent> { - public TrackerEventPersister(ReservedValueService reservedValueService) { - super(reservedValueService); + public TrackerEventPersister(ReservedValueService reservedValueService, DataSource dataSource) { + super(reservedValueService, dataSource); } @Override