Skip to content

Commit 055eaac

Browse files
Restore from backup - Git Bases Merge Conflicts
1 parent 7cdafff commit 055eaac

4 files changed

Lines changed: 44 additions & 113 deletions

File tree

apps/bfd-pipeline-idr/pipeline_stages.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@
3434
from model.idr_contract_pbp_contact import IdrContractPbpContact
3535
from model.idr_contract_pbp_number import IdrContractPbpNumber
3636
from model.idr_prior_auth import IdrPriorAuth
37+
from model.idr_prior_auth_item import IdrPriorAuthItem
3738
from parallel_executor import ParallelStagesExecutor, Stage
3839
from pipeline_utils import (
3940
extract_and_load,

apps/bfd-pipeline-idr/pipeline_utils.py

Lines changed: 43 additions & 89 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,5 @@
11
import time
22
from datetime import UTC, datetime, timedelta
3-
from typing import Any
43

54
import psycopg
65
from loguru import logger
@@ -16,10 +15,6 @@
1615
CLAIM_PROFESSIONAL_SS_TABLE,
1716
DEFAULT_MAX_DATE,
1817
DEFAULT_PARTITION,
19-
FISS_CLM_SOURCE,
20-
IDR_CLAIM_TABLE,
21-
MCS_CLM_SOURCE,
22-
PART_D_CLAIM_TYPE_CODES,
2318
PHASE_1_CUTOFF,
2419
)
2520
from extractor import PostgresExtractor, SnowflakeExtractor, Source
@@ -35,8 +30,10 @@
3530
from model.idr_beneficiary_ma_part_d_enrollment import IdrBeneficiaryMaPartDEnrollment
3631
from model.idr_beneficiary_ma_part_d_enrollment_rx import IdrBeneficiaryMaPartDEnrollmentRx
3732
from model.load_progress import LoadProgress
38-
from settings import BENEFICIARY_PART_D_PRUNE_BATCH_LIMIT, BENEFICIARY_PRUNE_BATCH_LIMIT
39-
33+
from settings import (
34+
BENEFICIARY_PART_D_PRUNE_BATCH_LIMIT,
35+
BENEFICIARY_PRUNE_BATCH_LIMIT,
36+
)
4037

4138
_SHARED_SYSTEM_CLAIM_ITEM_TABLES = {
4239
CLAIM_INSTITUTIONAL_SS_TABLE: CLAIM_INSTITUTIONAL_ITEM_SS_TABLE,
@@ -147,8 +144,8 @@ def prune_phase_1_ss_claims(
147144
item_table = _SHARED_SYSTEM_CLAIM_ITEM_TABLES.get(claim_table)
148145
if item_table is None:
149146
return True
150-
prune_cutoff_date = job_start - timedelta(days=PHASE_1_CUTOFF)
151147

148+
prune_cutoff_date = job_start - timedelta(days=PHASE_1_CUTOFF)
152149
logger.info("pruning phase 1 ss claims older than {}", prune_cutoff_date)
153150

154151
prune_query, params = stale_phase_1_claims_query(claim_table, item_table, prune_cutoff_date)
@@ -165,6 +162,10 @@ def prune_phase_1_ss_claims(
165162

166163
total_row_count += res.rowcount
167164
logger.info("pruned {} rows from {}", res.rowcount, item_table)
165+
166+
if res.rowcount == 0:
167+
logger.info("Total rows pruned from {}: {}", item_table, total_row_count)
168+
break
168169
return True
169170

170171

@@ -177,88 +178,41 @@ def prune_non_latest_non_part_d_ss_claims(
177178
item_table = _SHARED_SYSTEM_CLAIM_ITEM_TABLES.get(claim_table)
178179
if item_table is None:
179180
return True
180-
prune_cutoff_date = job_start - timedelta(days=PHASE_1_CUTOFF)
181-
part_d_codes = ",".join(str(code) for code in PART_D_CLAIM_TYPE_CODES)
182181

183-
non_latest_non_part_d_claim_filter = f"""
184-
clm.clm_ltst_clm_ind = 'N'
185-
AND clm.clm_type_cd NOT IN ({part_d_codes})
186-
AND clm.clm_idr_ld_dt < %s
187-
"""
182+
prune_cutoff_date = job_start - timedelta(days=PHASE_1_CUTOFF)
183+
logger.info("pruning non-latest non-Part-D ss claims older than {}", prune_cutoff_date)
188184

189-
logger.info("pruning non-latest non-Part-D ss claims older than {}", prune_cutoff_date)
185+
prune_query, params = non_latest_non_part_d_claims_query(claim_table, prune_cutoff_date)
190186

191-
if res.rowcount == 0:
192-
logger.info("Total rows pruned from {}: {}", item_table, total_row_count)
193-
break
194-
with psycopg.connect(get_connection_string(load_mode)) as conn:
195-
# Claim items can exist even when the non-latest parent claim was filtered
196-
# before final claim-table load, so use the source claim table.
197-
_prune_table_in_batches(
198-
conn,
199-
item_table,
200-
f"""
201-
DELETE FROM {item_table}
202-
WHERE (clm_uniq_id, bfd_row_id) IN (
203-
SELECT item.clm_uniq_id, item.bfd_row_id
204-
FROM {item_table} item
205-
JOIN {IDR_CLAIM_TABLE} clm ON clm.clm_uniq_id = item.clm_uniq_id
206-
WHERE {non_latest_non_part_d_claim_filter}
207-
LIMIT {PRUNE_BATCH_MAX_SIZE}
208-
)
209-
""",
210-
(prune_cutoff_date, PHASE_1_PRUNE_BATCH_LIMIT),
211-
)
212-
logger.info("pruned {} rows from {}", res.rowcount, item_table)
213-
if res.rowcount < PHASE_1_PRUNE_BATCH_LIMIT:
214-
break
215-
prune_query, params = non_latest_non_part_d_claim_items_query(item_table, prune_cutoff_date)
187+
total_row_counts = {
188+
item_table: 0,
189+
claim_table: 0,
190+
}
216191

217192
with psycopg.connect(get_connection_string(load_mode)) as conn:
218193
while True:
194+
claim_row_count = 0
219195
with conn.transaction():
220-
res = conn.execute(
221-
f"""
222-
DELETE FROM {item_table}
223-
WHERE (clm_uniq_id, bfd_row_id) IN ({prune_query})
224-
""", # type: ignore
225-
params,
226-
)
227-
logger.info("pruned {} rows from {}", res.rowcount, item_table)
228-
if res.rowcount < PHASE_1_PRUNE_BATCH_LIMIT:
229-
break
230-
231-
return True
232-
233-
234-
def prune_non_latest_non_part_d_ss_parent_claims(
235-
cls: type[T],
236-
load_mode: LoadMode,
237-
job_start: datetime,
238-
) -> bool:
239-
claim_table = cls.table()
240-
if claim_table not in _SHARED_SYSTEM_CLAIM_ITEM_TABLES:
241-
return True
242-
243-
prune_cutoff_date = job_start - timedelta(days=PHASE_1_CUTOFF)
196+
for target_table in [item_table, claim_table]:
197+
res = conn.execute(
198+
f"""DELETE FROM {target_table} WHERE clm_uniq_id IN ({prune_query})""", # type: ignore
199+
params,
200+
)
244201

245-
logger.info("pruning non-latest non-Part-D ss parent claims older than {}", prune_cutoff_date)
202+
total_row_counts[target_table] += res.rowcount
203+
logger.info("pruned {} rows from {}", res.rowcount, target_table)
246204

247-
prune_query, params = non_latest_non_part_d_parent_claims_query(claim_table, prune_cutoff_date)
205+
if target_table == claim_table:
206+
claim_row_count = res.rowcount
248207

249-
with psycopg.connect(get_connection_string(load_mode)) as conn:
250-
while True:
251-
with conn.transaction():
252-
res = conn.execute(
253-
f"""
254-
DELETE FROM {claim_table}
255-
WHERE clm_uniq_id IN ({prune_query})
256-
""", # type: ignore
257-
params,
258-
)
259-
logger.info("pruned {} rows from {}", res.rowcount, claim_table)
260-
if res.rowcount < PHASE_1_PRUNE_BATCH_LIMIT:
261-
break
208+
if claim_row_count == 0:
209+
for target_table in [item_table, claim_table]:
210+
logger.info(
211+
"Total rows pruned from {}: {}",
212+
target_table,
213+
total_row_counts[target_table],
214+
)
215+
break
262216

263217
return True
264218

@@ -276,7 +230,7 @@ def prune_bene_lis_cmbnd(
276230
f"""
277231
DELETE FROM {bene_table}
278232
WHERE (bene_sk, bene_cmbnd_deemd_efctv_dt, idr_trans_obslt_ts) IN (
279-
SELECT bene_sk, bene_cmbnd_deemd_efctv_dt, idr_trans_obslt_ts
233+
SELECT bene_sk, bene_cmbnd_deemd_efctv_dt, idr_trans_obslt_ts
280234
FROM {bene_table}
281235
WHERE idr_trans_obslt_ts < %s
282236
LIMIT %s
@@ -331,16 +285,16 @@ def prune_bene_ma_part_d_rx(
331285
res = conn.execute(
332286
f"""
333287
DELETE FROM {bene_table}
334-
WHERE (bene_sk,
335-
bene_cntrct_num,
336-
bene_pbp_num,
337-
bene_enrlmt_bgn_dt,
288+
WHERE (bene_sk,
289+
bene_cntrct_num,
290+
bene_pbp_num,
291+
bene_enrlmt_bgn_dt,
338292
bene_enrlmt_pdp_rx_info_bgn_dt
339293
) IN (
340-
SELECT bene_sk,
341-
bene_cntrct_num,
342-
bene_pbp_num,
343-
bene_enrlmt_bgn_dt,
294+
SELECT bene_sk,
295+
bene_cntrct_num,
296+
bene_pbp_num,
297+
bene_enrlmt_bgn_dt,
344298
bene_enrlmt_pdp_rx_info_bgn_dt
345299
FROM {bene_table}
346300
WHERE idr_trans_obslt_ts < %s

apps/bfd-pipeline-idr/test_pipeline.py

Lines changed: 0 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -322,27 +322,6 @@ def _do_test_pipeline(conn: Connection[DictRow], load_type: LoadType) -> None:
322322
rows = cur.fetchmany(1)
323323
assert rows[0]["clm_uniq_id"] == 849348853948
324324

325-
# Items for non-latest non-Part-D SS claims are pruned on incremental loads
326-
cur = conn.execute(
327-
"select * from idr.claim_item_institutional_ss where clm_uniq_id = 999999434800"
328-
)
329-
if load_type == LoadType.INITIAL:
330-
assert cur.rowcount == 1
331-
else:
332-
assert cur.rowcount == 0
333-
334-
cur = conn.execute("select * from idr.claim_item_professional_nch order by clm_uniq_id")
335-
assert cur.rowcount == 442
336-
rows = cur.fetchmany(1)
337-
assert rows[0]["clm_uniq_id"] == 119855147698
338-
339-
cur = conn.execute("select * from idr.claim_item_professional_ss order by clm_uniq_id")
340-
assert cur.rowcount == 1
341-
rows = cur.fetchmany(1)
342-
assert rows[0]["clm_uniq_id"] == 4991490559710
343-
344-
conn.commit()
345-
346325
# Test incremental loading logic involving 'source_load_events' if we're testing incremental
347326
# mode
348327
if load_type == LoadType.INCREMENTAL:

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

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -7,9 +7,6 @@
77
import gov.cms.bfd.server.ng.DbFilterBuilder;
88
import gov.cms.bfd.server.ng.DbFilterParam;
99
import gov.cms.bfd.server.ng.claim.model.*;
10-
import gov.cms.bfd.server.ng.claim.model.ClaimBase;
11-
import gov.cms.bfd.server.ng.claim.model.PriorAuthorization;
12-
import gov.cms.bfd.server.ng.claim.model.SystemType;
1310
import gov.cms.bfd.server.ng.input.ClaimSearchCriteria;
1411
import gov.cms.bfd.server.ng.log.QueryTelemetryUtil;
1512
import gov.cms.bfd.server.ng.util.MetricRecorder;

0 commit comments

Comments
 (0)