Skip to content

Commit 3fc6c39

Browse files
authored
Merge branch 'master' into BFD-4764
2 parents dfbbe5a + bb05a7f commit 3fc6c39

11 files changed

Lines changed: 143 additions & 22 deletions
Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,25 @@
1+
DROP VIEW idr.beneficiary_part_c_and_d_enrollment;
2+
CREATE VIEW idr.beneficiary_part_c_and_d_enrollment AS
3+
SELECT
4+
e.bene_sk,
5+
e.bene_enrlmt_pgm_type_cd,
6+
e.bene_enrlmt_bgn_dt,
7+
e.bene_enrlmt_end_dt,
8+
e.bene_cntrct_num,
9+
e.bene_pbp_num,
10+
e.cntrct_pbp_sk,
11+
e.bene_cvrg_type_cd,
12+
e.bene_enrlmt_emplr_sbsdy_sw,
13+
COALESCE(rx.bene_enrlmt_pdp_rx_info_bgn_dt, DATE '9999-12-31') AS bene_enrlmt_pdp_rx_info_bgn_dt,
14+
rx.bene_pdp_enrlmt_mmbr_id_num,
15+
rx.bene_pdp_enrlmt_grp_num,
16+
rx.bene_pdp_enrlmt_prcsr_num,
17+
rx.bene_pdp_enrlmt_bank_id_num
18+
FROM idr.beneficiary_ma_part_d_enrollment e
19+
LEFT JOIN idr.beneficiary_ma_part_d_enrollment_rx rx
20+
ON e.bene_sk = rx.bene_sk
21+
AND e.bene_enrlmt_bgn_dt = rx.bene_enrlmt_bgn_dt
22+
AND e.cntrct_pbp_sk = rx.cntrct_pbp_sk
23+
AND e.bene_enrlmt_pgm_type_cd in ('2', '3')
24+
AND rx.idr_trans_obslt_ts >= '9999-12-31'
25+
WHERE e.idr_trans_obslt_ts >= '9999-12-31';
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
CREATE INDEX ON idr.beneficiary_ma_part_d_enrollment(idr_trans_obslt_ts) WHERE idr_trans_obslt_ts < '9999-12-31';
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
CREATE INDEX ON idr.beneficiary_ma_part_d_enrollment_rx(idr_trans_obslt_ts) WHERE idr_trans_obslt_ts < '9999-12-31';

apps/bfd-pipeline-idr/model/idr_beneficiary_ma_part_d_enrollment.py

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -3,10 +3,7 @@
33

44
from pydantic import BeforeValidator
55

6-
from constants import (
7-
DEFAULT_MAX_DATE,
8-
IDR_BENE_MA_PART_D_TABLE,
9-
)
6+
from constants import IDR_BENE_MA_PART_D_TABLE
107
from load_partition import LoadPartition
118
from model.base_model import (
129
ALIAS_HSTRY,
@@ -76,7 +73,6 @@ def fetch_query(cls, partition: LoadPartition, start_time: datetime, source: Sou
7673
{deceased_bene_filter(hstry, start_time)}
7774
AND {hstry}.bene_sk = enrlmt.bene_sk
7875
)
79-
AND idr_trans_obslt_ts >= '{DEFAULT_MAX_DATE}'
8076
AND bene_enrlmt_pgm_type_cd != '~'
8177
{{ORDER_BY}}
8278
{{LIMIT}}

apps/bfd-pipeline-idr/model/idr_beneficiary_ma_part_d_enrollment_rx.py

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -3,10 +3,7 @@
33

44
from pydantic import BeforeValidator
55

6-
from constants import (
7-
DEFAULT_MAX_DATE,
8-
IDR_BENE_MA_PART_D_RX_TABLE,
9-
)
6+
from constants import IDR_BENE_MA_PART_D_RX_TABLE
107
from load_partition import LoadPartition
118
from model.base_model import (
129
ALIAS_HSTRY,
@@ -70,7 +67,6 @@ def fetch_query(cls, partition: LoadPartition, start_time: datetime, source: Sou
7067
{deceased_bene_filter(hstry, start_time)}
7168
AND {hstry}.bene_sk = enrlmt_rx.bene_sk
7269
)
73-
AND idr_trans_obslt_ts >= '{DEFAULT_MAX_DATE}'
7470
{{ORDER_BY}}
7571
{{LIMIT}}
7672
"""

apps/bfd-pipeline-idr/pipeline_stages.py

Lines changed: 19 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -36,7 +36,13 @@
3636
from model.idr_prior_auth import IdrPriorAuth
3737
from model.idr_prior_auth_item import IdrPriorAuthItem
3838
from parallel_executor import ParallelStagesExecutor, Stage
39-
from pipeline_utils import extract_and_load, prune_bene_lis_cmbnd, prune_phase_1_ss_claims
39+
from pipeline_utils import (
40+
extract_and_load,
41+
prune_bene_lis_cmbnd,
42+
prune_bene_ma_part_d,
43+
prune_bene_ma_part_d_rx,
44+
prune_phase_1_ss_claims,
45+
)
4046
from settings import enable_prior_auth_ingestion
4147

4248
type NodePartitionedModelInput = tuple[type[IdrBaseModel], LoadPartition | None]
@@ -107,7 +113,7 @@ async def start(self) -> bool:
107113
self._stage2_do_claims_and_benes_tbls(),
108114
self._stage3_do_parent_claims_tbls(),
109115
self._stage4_do_beneficiary(),
110-
self._stage5_do_phase_1_prune(),
116+
self._stage5_prune_obsolete_rows(),
111117
],
112118
)
113119
)
@@ -143,7 +149,7 @@ def _stage4_do_beneficiary(self) -> Stage[bool]:
143149
self._gen_partitioned_node_inputs(self._filter_tables(BENE_TABLES))
144150
)
145151

146-
def _stage5_do_phase_1_prune(self) -> Stage[bool]:
152+
def _stage5_prune_obsolete_rows(self) -> Stage[bool]:
147153
if self.load_type == LoadType.INITIAL:
148154
return
149155

@@ -152,6 +158,16 @@ def _stage5_do_phase_1_prune(self) -> Stage[bool]:
152158
self.load_mode,
153159
)
154160

161+
yield functools.partial(
162+
prune_bene_ma_part_d,
163+
self.load_mode,
164+
)
165+
166+
yield functools.partial(
167+
prune_bene_ma_part_d_rx,
168+
self.load_mode,
169+
)
170+
155171
for model in self._filter_tables(CLAIM_SS_TABLES):
156172
yield functools.partial(
157173
prune_phase_1_ss_claims,

apps/bfd-pipeline-idr/pipeline_utils.py

Lines changed: 68 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -26,8 +26,10 @@
2626
stale_phase_1_claims_query,
2727
)
2828
from model.idr_beneficiary_low_income_subsidy_cmbnd import IdrBeneficiaryLowIncomeSubsidyCmbnd
29+
from model.idr_beneficiary_ma_part_d_enrollment import IdrBeneficiaryMaPartDEnrollment
30+
from model.idr_beneficiary_ma_part_d_enrollment_rx import IdrBeneficiaryMaPartDEnrollmentRx
2931
from model.load_progress import LoadProgress
30-
from settings import BENEFICIARY_PRUNE_BATCH_LIMIT
32+
from settings import BENEFICIARY_PART_D_PRUNE_BATCH_LIMIT, BENEFICIARY_PRUNE_BATCH_LIMIT
3133

3234

3335
def get_progress(
@@ -189,3 +191,68 @@ def prune_bene_lis_cmbnd(
189191
break
190192

191193
return True
194+
195+
196+
def prune_bene_ma_part_d(
197+
load_mode: LoadMode,
198+
) -> bool:
199+
bene_table = IdrBeneficiaryMaPartDEnrollment.table()
200+
201+
logger.info("pruning obsolete part d beneficiaries", DEFAULT_MAX_DATE)
202+
203+
with psycopg.connect(get_connection_string(load_mode)) as conn, conn.transaction():
204+
while True:
205+
res = conn.execute(
206+
f"""
207+
DELETE FROM {bene_table}
208+
WHERE (bene_sk, bene_enrlmt_bgn_dt, bene_enrlmt_pgm_type_cd) IN (
209+
SELECT bene_sk, bene_enrlmt_bgn_dt, bene_enrlmt_pgm_type_cd
210+
FROM {bene_table}
211+
WHERE idr_trans_obslt_ts < %s
212+
LIMIT %s
213+
)
214+
""", # type: ignore
215+
(DEFAULT_MAX_DATE, BENEFICIARY_PART_D_PRUNE_BATCH_LIMIT),
216+
)
217+
logger.info("pruned {} rows from {}", res.rowcount, bene_table)
218+
if res.rowcount < BENEFICIARY_PART_D_PRUNE_BATCH_LIMIT:
219+
break
220+
221+
return True
222+
223+
224+
def prune_bene_ma_part_d_rx(
225+
load_mode: LoadMode,
226+
) -> bool:
227+
bene_table = IdrBeneficiaryMaPartDEnrollmentRx.table()
228+
229+
logger.info("pruning obsolete part d rx beneficiaries", DEFAULT_MAX_DATE)
230+
231+
with psycopg.connect(get_connection_string(load_mode)) as conn, conn.transaction():
232+
while True:
233+
res = conn.execute(
234+
f"""
235+
DELETE FROM {bene_table}
236+
WHERE (bene_sk,
237+
bene_cntrct_num,
238+
bene_pbp_num,
239+
bene_enrlmt_bgn_dt,
240+
bene_enrlmt_pdp_rx_info_bgn_dt
241+
) IN (
242+
SELECT bene_sk,
243+
bene_cntrct_num,
244+
bene_pbp_num,
245+
bene_enrlmt_bgn_dt,
246+
bene_enrlmt_pdp_rx_info_bgn_dt
247+
FROM {bene_table}
248+
WHERE idr_trans_obslt_ts < %s
249+
LIMIT %s
250+
)
251+
""", # type: ignore
252+
(DEFAULT_MAX_DATE, BENEFICIARY_PART_D_PRUNE_BATCH_LIMIT),
253+
)
254+
logger.info("pruned {} rows from {}", res.rowcount, bene_table)
255+
if res.rowcount < BENEFICIARY_PART_D_PRUNE_BATCH_LIMIT:
256+
break
257+
258+
return True

apps/bfd-pipeline-idr/settings.py

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -100,6 +100,11 @@ def enable_prior_auth_ingestion() -> bool:
100100
"""Number of minimum connections to hold in the pool concurrently per-batch for non-LOCAL loads.
101101
Defaults to 20."""
102102

103+
BENEFICIARY_PART_D_PRUNE_BATCH_LIMIT = int(
104+
getenv("IDR_BENEFICIARY_PART_D_PRUNE_BATCH_LIMIT", "100000")
105+
)
106+
"""Maximum rows to delete per prune statement for Part D beneficiary records."""
107+
103108
PHASE_1_PRUNE_BATCH_LIMIT = int(getenv("PHASE_1_PRUNE_BATCH_LIMIT", "10_000"))
104109
"""The maximum batch size for pruning old claims (phase 1 claims from shared systems) on
105110
INCREMENTAL loads. Defaults to 10000."""

apps/bfd-pipeline-idr/test_pipeline.py

Lines changed: 20 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -168,15 +168,27 @@ def _do_test_pipeline(conn: Connection[DictRow], load_type: LoadType) -> None:
168168
# only a future record exists for this contract
169169
assert rows[6]["cntrct_pbp_bgn_dt"].strftime("%Y-%m-%d") == "2026-12-01"
170170

171-
cur = conn.execute("select * from idr.beneficiary_ma_part_d_enrollment order by bene_sk")
172-
assert cur.rowcount == 3
173-
rows = cur.fetchmany(1)
174-
assert rows[0]["bene_sk"] == 353816020
171+
if load_type == LoadType.INITIAL:
172+
cur = conn.execute("select * from idr.beneficiary_ma_part_d_enrollment order by bene_sk")
173+
assert cur.rowcount == 4
174+
rows = cur.fetchmany(1)
175+
assert rows[0]["bene_sk"] == 353816020
176+
else:
177+
cur = conn.execute("select * from idr.beneficiary_ma_part_d_enrollment order by bene_sk")
178+
assert cur.rowcount == 3
179+
rows = cur.fetchmany(1)
180+
assert rows[0]["bene_sk"] == 353816020
175181

176-
cur = conn.execute("select * from idr.beneficiary_ma_part_d_enrollment_rx order by bene_sk")
177-
assert cur.rowcount == 2
178-
rows = cur.fetchmany(1)
179-
assert rows[0]["bene_sk"] == 353816020
182+
if load_type == LoadType.INITIAL:
183+
cur = conn.execute("select * from idr.beneficiary_ma_part_d_enrollment_rx order by bene_sk")
184+
assert cur.rowcount == 3
185+
rows = cur.fetchmany(1)
186+
assert rows[0]["bene_sk"] == 353816020
187+
else:
188+
cur = conn.execute("select * from idr.beneficiary_ma_part_d_enrollment_rx order by bene_sk")
189+
assert cur.rowcount == 2
190+
rows = cur.fetchmany(1)
191+
assert rows[0]["bene_sk"] == 353816020
180192

181193
cur = conn.execute("select * from idr.beneficiary_low_income_subsidy order by bene_sk")
182194
assert cur.rowcount == 2

apps/bfd-pipeline-idr/test_samples1/SYNTHETIC_BENE_MAPD_ENRLMT.csv

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,3 +2,4 @@ BENE_SK,CNTRCT_PBP_SK,IDR_LTST_TRANS_FLG,BENE_CNTRCT_NUM,BENE_PBP_NUM,BENE_CVRG_
22
547437476,993028062567,Y,S0001,003,11,2,1,2019-02-18,9999-12-31,2019-02-18T00:00:00.000000,2019-02-18T00:00:00.000000,2019-11-11T00:00:00.000000,9999-12-31T00:00:00.000000
33
441149422,408933975817,Y,G1234,002,3,1,~,2018-02-04,9999-12-31,2018-02-04T00:00:00.000000,2018-02-04T00:00:00.000000,2018-02-04T00:00:00.000000,9999-12-31T00:00:00.000000
44
353816020,761144926385,Y,H1234,003,11,3,Y,2020-09-03,9999-12-31,2020-09-03T00:00:00.000000,2020-09-03T00:00:00.000000,2020-09-03T00:00:00.000000,9999-12-31T00:00:00.000000
5+
353816021,993028062568,N,T0001,005,11,2,1,2020-06-01,2023-01-31,2020-06-01T00:00:00.000000,2020-06-01T00:00:00.000000,2020-06-01T00:00:00.000000,2023-01-15T00:00:00.000000

0 commit comments

Comments
 (0)