Skip to content

Commit 21ce29c

Browse files
BFD-4656: Optimize stream usage and process EOBs in parallel (#3084)
1 parent d00b503 commit 21ce29c

12 files changed

Lines changed: 85 additions & 74 deletions

File tree

apps/bfd-server-ng/src/main/java/gov/cms/bfd/server/ng/beneficiary/BeneficiaryRepository.java

Lines changed: 13 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -36,10 +36,10 @@ public List<BeneficiaryIdentity> getValidBeneficiaryIdentities(long beneXrefSk)
3636
return entityManager
3737
.createQuery(
3838
"""
39-
SELECT identity
40-
FROM BeneficiaryIdentity identity
41-
WHERE identity.id.xrefSk = :beneXrefSk
42-
""",
39+
SELECT identity
40+
FROM BeneficiaryIdentity identity
41+
WHERE identity.id.xrefSk = :beneXrefSk
42+
""",
4343
BeneficiaryIdentity.class)
4444
.setParameter("beneXrefSk", beneXrefSk)
4545
.getResultList();
@@ -96,10 +96,10 @@ public Optional<Long> getXrefSkFromBeneSk(long beneSk) {
9696
return entityManager
9797
.createQuery(
9898
"""
99-
SELECT bene.xrefSk
100-
FROM Beneficiary bene
101-
WHERE bene.beneSk = :beneSk
102-
""",
99+
SELECT bene.xrefSk
100+
FROM Beneficiary bene
101+
WHERE bene.beneSk = :beneSk
102+
""",
103103
Long.class)
104104
.setParameter("beneSk", beneSk)
105105
.getResultList()
@@ -118,10 +118,10 @@ public Optional<Long> getXrefSkFromMbi(String mbi) {
118118
return entityManager
119119
.createQuery(
120120
"""
121-
SELECT bene.xrefSk
122-
FROM Beneficiary bene
123-
WHERE bene.identifier.mbi = :mbi
124-
""",
121+
SELECT bene.xrefSk
122+
FROM Beneficiary bene
123+
WHERE bene.identifier.mbi = :mbi
124+
""",
125125
Long.class)
126126
.setParameter("mbi", mbi)
127127
.getResultList()
@@ -168,6 +168,7 @@ public PatientMatchResult searchPatientMatch(PatientMatch patientMatch) {
168168

169169
if (uniqueXrefs.size() == 1) {
170170
var matchedBene = benes.getFirst();
171+
LogUtil.logBeneSk(matchedBene.getBeneSk());
171172
var finalDetermination =
172173
new FinalDetermination(combinationIndex, matchedRecords.getFirst());
173174
return new PatientMatchResult(

apps/bfd-server-ng/src/main/java/gov/cms/bfd/server/ng/claim/ClaimAsyncService.java

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
import gov.cms.bfd.server.ng.DbFilterParam;
66
import gov.cms.bfd.server.ng.claim.model.*;
77
import gov.cms.bfd.server.ng.input.ClaimSearchCriteria;
8+
import gov.cms.bfd.server.ng.util.LogUtil;
89
import jakarta.persistence.EntityManager;
910
import jakarta.persistence.PersistenceContext;
1011
import java.util.*;
@@ -46,7 +47,7 @@ CompletableFuture<Optional<C>> findByIdInClaimType(
4647
.getResultList()
4748
.stream()
4849
.findFirst();
49-
50+
result.ifPresent(claim -> LogUtil.logBeneSk(claim.getBeneficiary().getBeneSk()));
5051
return CompletableFuture.completedFuture(result);
5152
}
5253

@@ -73,7 +74,9 @@ protected <T extends ClaimBase> CompletableFuture<List<T>> fetchClaims(
7374
DbFilterParam.withParams(entityManager.createQuery(jpql, claimClass), filters.params())
7475
.setParameter("beneSk", criteria.beneSk())
7576
.getResultList();
76-
77+
result.stream()
78+
.findFirst()
79+
.ifPresent(claim -> LogUtil.logBeneSk(claim.getBeneficiary().getBeneSk()));
7780
return CompletableFuture.completedFuture(result);
7881
}
7982

apps/bfd-server-ng/src/main/java/gov/cms/bfd/server/ng/claim/ClaimRepository.java

Lines changed: 25 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,6 @@
55
import gov.cms.bfd.server.ng.claim.model.*;
66
import gov.cms.bfd.server.ng.input.ClaimSearchCriteria;
77
import gov.cms.bfd.server.ng.input.DateTimeRange;
8-
import gov.cms.bfd.server.ng.util.LogUtil;
98
import io.micrometer.core.annotation.Timed;
109
import io.micrometer.core.aop.MeterTag;
1110
import java.util.*;
@@ -23,42 +22,42 @@ public class ClaimRepository {
2322

2423
private static final String CLAIM_PROFESSIONAL_SHARED_SYSTEMS =
2524
"""
26-
SELECT c
27-
FROM ClaimProfessionalSharedSystems c
28-
JOIN FETCH c.beneficiary b
29-
LEFT JOIN FETCH c.claimItems cl
30-
""";
25+
SELECT c
26+
FROM ClaimProfessionalSharedSystems c
27+
JOIN FETCH c.beneficiary b
28+
LEFT JOIN FETCH c.claimItems cl
29+
""";
3130

3231
private static final String CLAIM_PROFESSIONAL_NCH =
3332
"""
34-
SELECT c
35-
FROM ClaimProfessionalNch c
36-
JOIN FETCH c.beneficiary b
37-
JOIN FETCH c.claimItems cl
38-
""";
33+
SELECT c
34+
FROM ClaimProfessionalNch c
35+
JOIN FETCH c.beneficiary b
36+
JOIN FETCH c.claimItems cl
37+
""";
3938

4039
private static final String CLAIM_INSTITUTIONAL_SHARED_SYSTEMS =
4140
"""
42-
SELECT c
43-
FROM ClaimInstitutionalSharedSystems c
44-
JOIN FETCH c.beneficiary b
45-
LEFT JOIN FETCH c.claimItems cl
46-
""";
41+
SELECT c
42+
FROM ClaimInstitutionalSharedSystems c
43+
JOIN FETCH c.beneficiary b
44+
LEFT JOIN FETCH c.claimItems cl
45+
""";
4746

4847
private static final String CLAIM_INSTITUTIONAL_NCH =
4948
"""
50-
SELECT c
51-
FROM ClaimInstitutionalNch c
52-
JOIN FETCH c.beneficiary b
53-
JOIN FETCH c.claimItems cl
54-
""";
49+
SELECT c
50+
FROM ClaimInstitutionalNch c
51+
JOIN FETCH c.beneficiary b
52+
JOIN FETCH c.claimItems cl
53+
""";
5554

5655
private static final String CLAIM_RX =
5756
"""
58-
SELECT c
59-
FROM ClaimRx c
60-
JOIN FETCH c.beneficiary b
61-
""";
57+
SELECT c
58+
FROM ClaimRx c
59+
JOIN FETCH c.beneficiary b
60+
""";
6261

6362
private static final List<ClaimTypeDefinition> ALL_CLAIM_TYPES =
6463
List.of(
@@ -152,11 +151,7 @@ public Optional<ClaimBase> findById(
152151
institutionalNchClaims.join().ifPresent(allClaims::add);
153152
rxClaims.join().ifPresent(allClaims::add);
154153

155-
var result = allClaims.stream().findFirst();
156-
157-
result.ifPresent(claim -> LogUtil.logBeneSk(claim.getBeneficiary().getBeneSk()));
158-
159-
return result;
154+
return allClaims.stream().findFirst();
160155
}
161156

162157
/**

apps/bfd-server-ng/src/main/java/gov/cms/bfd/server/ng/coverage/CoverageHandler.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -7,8 +7,8 @@
77
import gov.cms.bfd.server.ng.loadprogress.LoadProgressRepository;
88
import gov.cms.bfd.server.ng.util.FhirUtil;
99
import java.util.Arrays;
10-
import java.util.List;
1110
import java.util.Optional;
11+
import java.util.stream.Stream;
1212
import lombok.RequiredArgsConstructor;
1313
import org.hl7.fhir.r4.model.Bundle;
1414
import org.hl7.fhir.r4.model.Coverage;
@@ -75,7 +75,7 @@ public Bundle searchByBeneficiary(Long beneSk, DateTimeRange lastUpdated) {
7575
.searchBeneficiaryWithCoverage(beneSk, lastUpdated)
7676
.filter(b -> !b.isMergedBeneficiary());
7777
if (beneficiaryOpt.isEmpty()) {
78-
return FhirUtil.bundleOrDefault(List.of(), loadProgressRepository::lastUpdated);
78+
return FhirUtil.bundleOrDefault(Stream.of(), loadProgressRepository::lastUpdated);
7979
}
8080
var beneficiary = beneficiaryOpt.get();
8181
var coverages =

apps/bfd-server-ng/src/main/java/gov/cms/bfd/server/ng/eob/EobHandler.java

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -96,14 +96,18 @@ public Bundle searchByBene(ClaimSearchCriteria criteria, SamhsaFilterMode samhsa
9696
return claim.toFhir(securityStatus);
9797
});
9898

99-
return FhirUtil.bundleOrDefault(filteredClaims.toList(), loadProgressRepository::lastUpdated);
99+
return FhirUtil.bundleOrDefault(filteredClaims, loadProgressRepository::lastUpdated);
100100
}
101101

102102
private Stream<? extends ClaimBase> filterSamhsaClaims(
103103
List<? extends ClaimBase> claims,
104104
SamhsaFilterMode samhsaFilterMode,
105105
ClaimSearchCriteria claimSearchCriteria) {
106-
var claimStream = claims.stream().sorted(Comparator.comparing(ClaimBase::getClaimUniqueId));
106+
// Process claims in parallel
107+
// Note: DO NOT call toList() until the very end as materializing the list multiple times could
108+
// negatively impact perf.
109+
var claimStream =
110+
claims.parallelStream().sorted(Comparator.comparing(ClaimBase::getClaimUniqueId));
107111
var filteredClaimStream =
108112
switch (samhsaFilterMode) {
109113
case INCLUDE -> claimStream;

apps/bfd-server-ng/src/main/java/gov/cms/bfd/server/ng/input/FhirInputConverter.java

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -101,7 +101,7 @@ public static Optional<Integer> toIntOptional(@Nullable NumberParam numberParam)
101101
}
102102
try {
103103
return Optional.of(numberParam.getValue().intValueExact());
104-
} catch (ArithmeticException ex) {
104+
} catch (ArithmeticException _) {
105105
throw new InvalidRequestException("Numeric input was not in a valid format");
106106
}
107107
}
@@ -114,8 +114,12 @@ public static Optional<Integer> toIntOptional(@Nullable NumberParam numberParam)
114114
* @return ID
115115
*/
116116
public static Long toLong(@Nullable ReferenceParam reference, String validResourceType) {
117-
if (reference == null || reference.getIdPartAsLong() == null) {
118-
throw new InvalidRequestException("Reference is missing");
117+
try {
118+
if (reference == null || reference.getIdPartAsLong() == null) {
119+
throw new InvalidRequestException("Reference is missing");
120+
}
121+
} catch (NumberFormatException _) {
122+
throw new InvalidRequestException("Reference is not a valid number");
119123
}
120124
var resourceType = reference.getResourceType();
121125
if (!StringUtils.isBlank(resourceType) && !resourceType.equals(validResourceType)) {

apps/bfd-server-ng/src/main/java/gov/cms/bfd/server/ng/log/LogStreamAuditLogger.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,7 @@ public void log(PatientMatchAuditRecord auditRecord) {
2727
LOGGER
2828
.atInfo()
2929
.setMessage(PATIENT_MATCH_REQUESTED)
30-
.addKeyValue("logType", "patientMatchAudit")
30+
.addKeyValue(LOG_TYPE, "patientMatchAudit")
3131
.addKeyValue(logKey(AUDIT_PREFIX, MATCHED_BENE_SK), matchedBeneSk.get().toString())
3232
.addKeyValue(
3333
logKey(AUDIT_PREFIX, BENE_SKS_FOUND),

apps/bfd-server-ng/src/main/java/gov/cms/bfd/server/ng/patient/PatientHandler.java

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,6 @@
1717
import gov.cms.bfd.server.ng.util.*;
1818
import java.time.Instant;
1919
import java.util.Arrays;
20-
import java.util.List;
2120
import java.util.Optional;
2221
import java.util.UUID;
2322
import java.util.stream.Stream;
@@ -126,7 +125,7 @@ public Bundle searchByBeneficiaryC4DIC(Long beneSk) {
126125
.searchBeneficiaryWithCoverage(beneSk, new DateTimeRange())
127126
.filter(b -> !b.isMergedBeneficiary());
128127
if (beneficiaryOpt.isEmpty()) {
129-
return FhirUtil.bundleOrDefault(List.of(), loadProgressRepository::lastUpdated);
128+
return FhirUtil.bundleOrDefault(Stream.of(), loadProgressRepository::lastUpdated);
130129
}
131130
var beneficiary = beneficiaryOpt.get();
132131

apps/bfd-server-ng/src/main/java/gov/cms/bfd/server/ng/util/FhirUtil.java

Lines changed: 5 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,6 @@
11
package gov.cms.bfd.server.ng.util;
22

33
import java.time.ZonedDateTime;
4-
import java.util.List;
54
import java.util.Optional;
65
import java.util.function.Supplier;
76
import java.util.regex.Pattern;
@@ -58,24 +57,13 @@ public static String getHcpcsSystem(String code) {
5857
* @return bundle
5958
*/
6059
public static Bundle bundleOrDefault(
61-
Stream<Resource> resources, Supplier<ZonedDateTime> batchLastUpdated) {
62-
return bundleOrDefault(resources.toList(), batchLastUpdated);
63-
}
60+
Stream<? extends Resource> resources, Supplier<ZonedDateTime> batchLastUpdated) {
61+
var bundle = getBundle(resources);
6462

65-
/**
66-
* Creates a bundle from the resource, returning a default bundle with lastUpdated populated if
67-
* empty.
68-
*
69-
* @param resources resources
70-
* @param batchLastUpdated last updated
71-
* @return bundle
72-
*/
73-
public static Bundle bundleOrDefault(
74-
List<? extends Resource> resources, Supplier<ZonedDateTime> batchLastUpdated) {
75-
if (resources.isEmpty()) {
63+
if (bundle.getEntry().isEmpty()) {
7664
return defaultBundle(batchLastUpdated);
7765
}
78-
return getBundle(resources.stream());
66+
return bundle;
7967
}
8068

8169
/**
@@ -89,7 +77,7 @@ public static Bundle bundleOrDefault(
8977
public static Bundle bundleOrDefault(
9078
Optional<Resource> resource, Supplier<ZonedDateTime> batchLastUpdated) {
9179
return resource
92-
.map(value -> bundleOrDefault(List.of(value), batchLastUpdated))
80+
.map(value -> bundleOrDefault(Stream.of(value), batchLastUpdated))
9381
.orElseGet(() -> defaultBundle(batchLastUpdated));
9482
}
9583

apps/bfd-server-ng/src/main/java/gov/cms/bfd/server/ng/util/LogUtil.java

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,7 @@
11
package gov.cms.bfd.server.ng.util;
22

3+
import static gov.cms.bfd.server.ng.util.LoggerConstants.LOG_TYPE;
4+
35
import org.slf4j.LoggerFactory;
46

57
/** Utility class used for logging. */
@@ -20,8 +22,9 @@ private LogUtil() {
2022
public static void logBeneSk(Long beneSk) {
2123
LOGGER
2224
.atInfo()
23-
.setMessage(LoggerConstants.BENE_SK_REQUESTED)
24-
.addKeyValue("bene_sk", beneSk)
25+
.setMessage(LoggerConstants.BENE_SK_FOUND)
26+
.addKeyValue(LOG_TYPE, "beneFound")
27+
.addKeyValue("beneSk", beneSk)
2528
.log();
2629
}
2730
}

0 commit comments

Comments
 (0)