cdc with python and merge on bigquery

Viewed 78

I have written a pipeline using Apache beam and Google dataflow that sends changes from a MongoDB to bq. I have a bigquery log table like ...

table operation type timestamp
[all columns] [insert / update / delete / replace] timestamp

and a "normal" table without the operation and timestamp column. My goal is to merge the src table (log) and the target table. The problem is as following, when the second last entry to a field is not null and the last one is, how can I check this in the merge statement? For example in other databases you can do something like

create function get_sec_last_value(id) as (
  (
    select as struct
      *
    from (
      select
        *,
        row_number() over(order by timestamp desc) as number
      from table
      where id = id
    ) where number = 2
  )
);

merge target trg
using source as src
on trg.id = src.id

...
update set id = case 
    when (get_sec_last_value(src.id).id is not null and src.id is not null) or (get_sec_last_value(src.id).id is null and src.id is not null) then src.id
    when (get_sec_last_value(src.id).id is not null and src.id is null) or (get_sec_last_value(src.id).id is null and src.id is null) then null
    end
...

Has anybody faced the same problem or has an idea how to solve it?

Thanks in advance

0 Answers
Related