Options
- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
07-29-2024 05:33 PM
Thanks @xorbix_rshiva for your reply,
There was deduplication logic in my code. Could I use _source_cdc_time for the watermark in this case?
def merge_stream(microBatchDF, i):
microBatchDF.createOrReplaceTempView("vw_delta")
microBatchDF._jdf.sparkSession().sql("""
MERGE INTO customer as tg
USING (
SELECT * FROM (
SELECT
key.customerid,
key.customerkey,
value.after.customername,
--meta data
value.op,
current_timestamp() as _ingested_time,
to_timestamp(value.source.ts_ms/1000) as _source_cdc_time,
to_timestamp(value.ts_ms/1000) as _kafka_created_time,
landing_time as _landing_time,
row_number() over(PARTITION BY key.customerid, key.cusomterkey order by value.source.ts_ms desc) as rank
FROM vw_delta
WHERE value.op IS NOT NULL
) as t
WHERE rank = 1
) as src
ON tg.customerid= src.customerid
AND tg.customerkey= src.customerkey
WHEN MATCHED AND src.op = 'd' THEN DELETE
WHEN MATCHED AND src.op != 'd' THEN UPDATE SET *
WHEN NOT MATCHED AND src.op != 'd' THEN INSERT *
""")