Skip to content

Commit 9c2319a

Browse files
BFD-4205: Refactor xref and kill credit handling (#2782)
1 parent 790bb93 commit 9c2319a

67 files changed

Lines changed: 993 additions & 706 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

apps/bfd-model-idr/generator_util.py

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -127,6 +127,9 @@ def generate_bene_xref(self, new_bene_sk, old_bene_sk):
127127
src_rec_ctre_ts = self.fake.date_time_between_dates(
128128
efctv_ts, datetime.datetime.now() - datetime.timedelta(days=1)
129129
)
130+
src_rec_updt_ts = self.fake.date_time_between_dates(
131+
src_rec_ctre_ts, datetime.datetime.now() - datetime.timedelta(days=1)
132+
)
130133
insrt_ts = self.fake.date_time_between_dates(
131134
efctv_ts, datetime.datetime.now() - datetime.timedelta(days=1)
132135
)
@@ -140,6 +143,7 @@ def generate_bene_xref(self, new_bene_sk, old_bene_sk):
140143
"BENE_HICN_NUM": bene_hicn_num,
141144
"BENE_KILL_CRED_CD": str(kill_cred_cd),
142145
"SRC_REC_CRTE_TS": str(src_rec_ctre_ts),
146+
"SRC_REC_UPDT_TS": str(src_rec_updt_ts),
143147
"IDR_TRANS_EFCTV_TS": str(efctv_ts),
144148
"IDR_INSRT_TS": str(insrt_ts),
145149
"IDR_UPDT_TS": str(updt_ts),
@@ -456,6 +460,7 @@ def save_output_files(self):
456460
[
457461
"BENE_SK",
458462
"BENE_XREF_EFCTV_SK",
463+
"BENE_XREF_SK",
459464
"BENE_MBI_ID",
460465
"BENE_LAST_NAME",
461466
"BENE_1ST_NAME",

apps/bfd-model-idr/patient_generator.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -146,6 +146,7 @@
146146

147147
patient["BENE_SK"] = str(pt_bene_sk)
148148
patient["BENE_XREF_EFCTV_SK"] = str(pt_bene_sk)
149+
patient["BENE_XREF_SK"] = patient["BENE_XREF_EFCTV_SK"]
149150
generator.used_bene_sk.append(pt_bene_sk)
150151

151152
if pd.notna(row.get("BENE_MBI_ID")):
@@ -186,6 +187,7 @@
186187
pt_bene_sk = generator.gen_bene_sk()
187188
patient["BENE_SK"] = str(pt_bene_sk)
188189
patient["BENE_XREF_EFCTV_SK"] = str(pt_bene_sk)
190+
patient["BENE_XREF_SK"] = patient["BENE_XREF_EFCTV_SK"]
189191
generator.used_bene_sk.append(pt_bene_sk)
190192

191193
num_mbis = random.choices([1, 2, 3, 4], weights=[0.8, 0.14, 0.05, 0.01])[0]

apps/bfd-pipeline/bfd-pipeline-idr/bfd.sql

Lines changed: 21 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -2,10 +2,12 @@ DROP SCHEMA IF EXISTS idr CASCADE;
22

33
CREATE SCHEMA idr;
44
CREATE TABLE idr.beneficiary(
5+
-- columns from V2_MDCR_BENE_HSTRY
56
bene_sk BIGINT NOT NULL,
67
bene_xref_efctv_sk BIGINT NOT NULL,
78
bene_xref_efctv_sk_computed BIGINT NOT NULL GENERATED ALWAYS
8-
AS (CASE WHEN bene_xref_efctv_sk = 0 THEN bene_sk ELSE bene_xref_efctv_sk END) STORED,
9+
-- if kill credit is set to 1, the merge is invalid, and we should act like it's not present
10+
AS (CASE WHEN bene_xref_efctv_sk = 0 OR bene_kill_cred_cd = '1' THEN bene_sk ELSE bene_xref_efctv_sk END) STORED,
911
bene_mbi_id VARCHAR(11) NOT NULL,
1012
bene_1st_name VARCHAR(30) NOT NULL,
1113
bene_midl_name VARCHAR(15) NOT NULL,
@@ -28,8 +30,13 @@ CREATE TABLE idr.beneficiary(
2830
idr_ltst_trans_flg VARCHAR(1) NOT NULL,
2931
idr_trans_efctv_ts TIMESTAMPTZ NOT NULL,
3032
idr_trans_obslt_ts TIMESTAMPTZ NOT NULL,
31-
idr_insrt_ts TIMESTAMPTZ NOT NULL,
32-
idr_updt_ts TIMESTAMPTZ NOT NULL,
33+
idr_insrt_ts_bene TIMESTAMPTZ NOT NULL,
34+
idr_updt_ts_bene TIMESTAMPTZ NOT NULL,
35+
-- columns from V2_MDCR_BENE_XREF
36+
bene_kill_cred_cd VARCHAR(1) NOT NULL,
37+
src_rec_updt_ts TIMESTAMPTZ NOT NULL,
38+
idr_insrt_ts_xref TIMESTAMPTZ NOT NULL,
39+
idr_updt_ts_xref TIMESTAMPTZ NOT NULL,
3340
bfd_created_ts TIMESTAMPTZ NOT NULL,
3441
bfd_updated_ts TIMESTAMPTZ NOT NULL,
3542
PRIMARY KEY(bene_sk, idr_trans_efctv_ts)
@@ -129,21 +136,6 @@ CREATE TABLE idr.beneficiary_election_period_usage (
129136
PRIMARY KEY(bene_sk, cntrct_pbp_sk, bene_enrlmt_efctv_dt)
130137
);
131138

132-
CREATE TABLE idr.beneficiary_xref (
133-
bene_sk BIGINT NOT NULL,
134-
bene_xref_sk BIGINT NOT NULL,
135-
bene_hicn_num VARCHAR(11) NOT NULL,
136-
bene_kill_cred_cd VARCHAR(1) NOT NULL,
137-
idr_insrt_ts TIMESTAMPTZ NOT NULL,
138-
idr_updt_ts TIMESTAMPTZ NOT NULL,
139-
src_rec_crte_ts TIMESTAMPTZ NOT NULL,
140-
idr_trans_efctv_ts TIMESTAMPTZ NOT NULL,
141-
idr_trans_obslt_ts TIMESTAMPTZ NOT NULL,
142-
bfd_created_ts TIMESTAMPTZ NOT NULL,
143-
bfd_updated_ts TIMESTAMPTZ NOT NULL,
144-
PRIMARY KEY(bene_sk, bene_hicn_num, src_rec_crte_ts)
145-
);
146-
147139
CREATE TABLE idr.contract_pbp_number (
148140
cntrct_pbp_sk BIGINT NOT NULL PRIMARY KEY,
149141
cntrct_drug_plan_ind_cd VARCHAR(1) NOT NULL,
@@ -447,6 +439,17 @@ CREATE VIEW idr.beneficiary_entitlement_reason_current AS
447439
SELECT * FROM idr.beneficiary_entitlement_reason
448440
WHERE idr_ltst_trans_flg = 'Y' AND bene_rng_bgn_dt <= NOW() - INTERVAL '12 hours' AND bene_rng_end_dt >= NOW() - INTERVAL '12 hours';
449441

442+
CREATE VIEW idr.beneficiary_identity AS
443+
SELECT DISTINCT
444+
bene.bene_sk,
445+
bene.bene_xref_efctv_sk_computed,
446+
bene.bene_mbi_id,
447+
bene_mbi.bene_mbi_efctv_dt,
448+
bene_mbi.bene_mbi_obslt_dt
449+
FROM idr.beneficiary bene
450+
LEFT JOIN idr.beneficiary_mbi_id bene_mbi
451+
ON bene.bene_mbi_id = bene_mbi.bene_mbi_id;
452+
450453
CREATE OR REPLACE FUNCTION idr.refresh_overshare_mbis()
451454
RETURNS VOID AS $$
452455
DECLARE comment_sql TEXT;

apps/bfd-pipeline/bfd-pipeline-idr/mock-idr.sql

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,8 @@ DROP SCHEMA IF EXISTS cms_vdm_view_mdcr_prd CASCADE;
33
CREATE SCHEMA cms_vdm_view_mdcr_prd;
44

55
CREATE TABLE cms_vdm_view_mdcr_prd.v2_mdcr_bene_hstry (
6-
bene_sk BIGINT NOT NULL,
6+
bene_sk BIGINT NOT NULL,
7+
bene_xref_sk BIGINT NOT NULL,
78
bene_xref_efctv_sk BIGINT NOT NULL,
89
bene_mbi_id VARCHAR(11),
910
bene_1st_name VARCHAR(30),
@@ -122,6 +123,7 @@ CREATE TABLE cms_vdm_view_mdcr_prd.v2_mdcr_bene_xref (
122123
idr_insrt_ts TIMESTAMPTZ NOT NULL,
123124
idr_updt_ts TIMESTAMPTZ,
124125
src_rec_crte_ts TIMESTAMPTZ NOT NULL,
126+
src_rec_updt_ts TIMESTAMPTZ NOT NULL,
125127
PRIMARY KEY(bene_sk, bene_hicn_num, src_rec_crte_ts)
126128
);
127129

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

Lines changed: 74 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -59,6 +59,17 @@ def transform_null_int(value: int | None) -> float:
5959
DERIVED = "derived"
6060
COLUMN_MAP = "column_map"
6161

62+
ALIAS_CLM = "clm"
63+
ALIAS_DCMTN = "dcmtn"
64+
ALIAS_SGNTR = "sgntr"
65+
ALIAS_LINE = "line"
66+
ALIAS_PROCEDURE = "prod"
67+
ALIAS_INSTNL = "instnl"
68+
ALIAS_PRFNL = "prfnl"
69+
ALIAS_VAL = "val"
70+
ALIAS_HSTRY = "hstry"
71+
ALIAS_XREF = "xref"
72+
6273

6374
class IdrBaseModel(BaseModel):
6475
@staticmethod
@@ -160,7 +171,8 @@ def insert_keys(cls) -> list[str]:
160171

161172

162173
class IdrBeneficiary(IdrBaseModel):
163-
bene_sk: Annotated[int, {PRIMARY_KEY: True}]
174+
# columns from V2_MDCR_BENE_HSTRY
175+
bene_sk: Annotated[int, {PRIMARY_KEY: True, ALIAS: ALIAS_HSTRY}]
164176
bene_xref_efctv_sk: int
165177
bene_mbi_id: Annotated[str, BeforeValidator(transform_null_string)]
166178
bene_1st_name: str
@@ -184,9 +196,26 @@ class IdrBeneficiary(IdrBaseModel):
184196
idr_ltst_trans_flg: Annotated[str, BeforeValidator(transform_null_string)]
185197
idr_trans_efctv_ts: Annotated[datetime, {PRIMARY_KEY: True}]
186198
idr_trans_obslt_ts: datetime
187-
idr_insrt_ts: Annotated[datetime, {BATCH_TIMESTAMP: True}]
188-
idr_updt_ts: Annotated[
189-
datetime, {UPDATE_TIMESTAMP: True}, BeforeValidator(transform_null_date_to_min)
199+
idr_insrt_ts_bene: Annotated[
200+
datetime, {BATCH_TIMESTAMP: True, ALIAS: ALIAS_HSTRY, COLUMN_MAP: "idr_insrt_ts"}
201+
]
202+
idr_updt_ts_bene: Annotated[
203+
datetime,
204+
{UPDATE_TIMESTAMP: True, ALIAS: ALIAS_HSTRY, COLUMN_MAP: "idr_updt_ts"},
205+
BeforeValidator(transform_null_date_to_min),
206+
]
207+
# columns from V2_MDCR_BENE_XREF
208+
bene_kill_cred_cd: Annotated[str, BeforeValidator(transform_default_string)]
209+
src_rec_updt_ts: Annotated[datetime, BeforeValidator(transform_null_date_to_min)]
210+
idr_insrt_ts_xref: Annotated[
211+
datetime,
212+
{BATCH_TIMESTAMP: True, ALIAS: ALIAS_XREF, COLUMN_MAP: "idr_insrt_ts"},
213+
BeforeValidator(transform_null_date_to_min),
214+
]
215+
idr_updt_ts_xref: Annotated[
216+
datetime,
217+
{UPDATE_TIMESTAMP: True, ALIAS: ALIAS_XREF, COLUMN_MAP: "idr_updt_ts"},
218+
BeforeValidator(transform_null_date_to_min),
190219
]
191220

192221
@staticmethod
@@ -199,11 +228,47 @@ def computed_keys() -> list[str]:
199228

200229
@staticmethod
201230
def _current_fetch_query(start_time: datetime) -> str: # noqa: ARG004
202-
return """
203-
SELECT {COLUMNS}
204-
FROM cms_vdm_view_mdcr_prd.v2_mdcr_bene_hstry
205-
{WHERE_CLAUSE}
206-
{ORDER_BY}
231+
hstry = ALIAS_HSTRY
232+
xref = ALIAS_XREF
233+
# There can be multiple xref records for the same bene_sk/bene_ref_sk combo
234+
# so we need to find the most recent one based on src_rec_updt_ts.
235+
236+
# Unlike idr_updt_ts, src_rec_updt_ts will be set to the created timestamp
237+
# if no update has been applied. Therefore, we can just check the updated timestamp
238+
# without caring about the created timestamp.
239+
240+
# There can also be duplicate values with the same idr_insrt_ts, so we have to rely on
241+
# src_rec_insrt_ts/src_rec_updt_ts for this.
242+
return f"""
243+
WITH ordered_xref AS (
244+
SELECT bene_sk, bene_xref_sk, ROW_NUMBER() OVER (
245+
PARTITION BY bene_sk, bene_xref_sk
246+
ORDER BY src_rec_updt_ts DESC
247+
) AS row_order
248+
FROM cms_vdm_view_mdcr_prd.v2_mdcr_bene_xref
249+
),
250+
current_xref AS (
251+
SELECT
252+
ox.bene_sk,
253+
ox.bene_xref_sk,
254+
bx.bene_kill_cred_cd,
255+
bx.src_rec_updt_ts,
256+
bx.idr_insrt_ts,
257+
bx.idr_updt_ts
258+
FROM ordered_xref ox
259+
JOIN cms_vdm_view_mdcr_prd.v2_mdcr_bene_xref bx
260+
ON bx.bene_sk = ox.bene_sk AND bx.bene_xref_sk = ox.bene_xref_sk
261+
WHERE ox.row_order = 1
262+
)
263+
SELECT {{COLUMNS}}
264+
FROM cms_vdm_view_mdcr_prd.v2_mdcr_bene_hstry {hstry}
265+
-- NOTE: the join condition is intentionally inverted here
266+
-- In the xref table, the bene_sk and bene_xref_sk fields are mirrored
267+
LEFT JOIN current_xref {xref}
268+
ON {xref}.bene_sk = {hstry}.bene_xref_sk
269+
AND {xref}.bene_xref_sk = {hstry}.bene_sk
270+
{{WHERE_CLAUSE}}
271+
{{ORDER_BY}}
207272
"""
208273

209274

@@ -344,33 +409,6 @@ def _current_fetch_query(start_time: datetime) -> str: # noqa: ARG004
344409
"""
345410

346411

347-
class IdrBeneficiaryXref(IdrBaseModel):
348-
bene_hicn_num: Annotated[str, {PRIMARY_KEY: True}]
349-
bene_sk: Annotated[int, {PRIMARY_KEY: True}]
350-
bene_xref_sk: int
351-
bene_kill_cred_cd: Annotated[str, BeforeValidator(transform_default_string)]
352-
idr_trans_efctv_ts: datetime
353-
idr_trans_obslt_ts: datetime
354-
idr_insrt_ts: datetime
355-
idr_updt_ts: Annotated[
356-
datetime, {UPDATE_TIMESTAMP: True}, BeforeValidator(transform_null_date_to_min)
357-
]
358-
src_rec_crte_ts: Annotated[datetime, {PRIMARY_KEY: True, BATCH_TIMESTAMP: True}]
359-
360-
@staticmethod
361-
def table() -> str:
362-
return "idr.beneficiary_xref"
363-
364-
@staticmethod
365-
def _current_fetch_query(start_time: datetime) -> str: # noqa: ARG004
366-
return """
367-
SELECT {COLUMNS}
368-
FROM cms_vdm_view_mdcr_prd.v2_mdcr_bene_xref
369-
{WHERE_CLAUSE}
370-
{ORDER_BY}
371-
"""
372-
373-
374412
class IdrElectionPeriodUsage(IdrBaseModel):
375413
bene_sk: Annotated[int, {PRIMARY_KEY: True}]
376414
cntrct_pbp_sk: Annotated[int, {PRIMARY_KEY: True}]
@@ -422,16 +460,6 @@ def _current_fetch_query(start_time: datetime) -> str: # noqa: ARG004
422460
"""
423461

424462

425-
ALIAS_CLM = "clm"
426-
ALIAS_DCMTN = "dcmtn"
427-
ALIAS_SGNTR = "sgntr"
428-
ALIAS_LINE = "line"
429-
ALIAS_PROCEDURE = "prod"
430-
ALIAS_INSTNL = "instnl"
431-
ALIAS_PRFNL = "prfnl"
432-
ALIAS_VAL = "val"
433-
434-
435463
def claim_type_clause(start_time: datetime) -> str: # noqa: ARG001
436464
latest_claims_env = "IDR_LATEST_CLAIMS"
437465
if latest_claims_env in os.environ and os.environ[latest_claims_env] in ("1", "true"):

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

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,6 @@
1717
IdrBeneficiaryMbiId,
1818
IdrBeneficiaryStatus,
1919
IdrBeneficiaryThirdParty,
20-
IdrBeneficiaryXref,
2120
IdrClaim,
2221
IdrClaimAnsiSignature,
2322
IdrClaimDateSignature,
@@ -142,7 +141,6 @@ def run_pipeline(data_extractor: Extractor, connection_string: str) -> None:
142141
IdrBeneficiaryThirdParty,
143142
IdrBeneficiaryEntitlement,
144143
IdrBeneficiaryEntitlementReason,
145-
IdrBeneficiaryXref,
146144
# Ignore for now, we'll likely source these elsewhere when we load contract data
147145
# IdrContractPbpNumber,
148146
# IdrElectionPeriodUsage,

apps/bfd-pipeline/bfd-pipeline-idr/test_pipeline.py

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -75,10 +75,12 @@ def test_pipeline(self, psql_url: str) -> None:
7575
assert rows[0]["bene_mbi_id"] == "1S000000000"
7676
assert rows[1]["bene_mbi_id"] == "5B88XK5JN88"
7777

78-
cur = conn.execute("select * from idr.beneficiary_xref")
78+
cur = conn.execute(
79+
"select * from idr.beneficiary where bene_kill_cred_cd != '' order by bene_sk"
80+
)
7981
assert cur.rowcount == 6
8082
rows = cur.fetchmany(1)
81-
assert rows[0]["bene_sk"] == 454323619
83+
assert rows[0]["bene_sk"] == 174441863
8284

8385
cur = conn.execute("select * from idr.beneficiary_third_party order by bene_sk")
8486
assert cur.rowcount == 4

0 commit comments

Comments
 (0)