1
votes

I like to share with you something I managed to make for Amazon Redshift which is dynamic merge statement for SCD(Slowly Changing Dimension) Type 2.

What this statement assumes exist beforehand:

  • Two schemas in database - dbimports and repo. dbimports schema is used as staging area and repo will be the target for SCD type 2 dimensions and facts.
  • That tables in repo and dbimports schema have the same names along with same column names coming from source data.
  • Tables in repo schema will have preset primary keys.
  • Tables in repo schema will have three additional columns - scd_key,is_active,inserted_date . scd_key is identity column and can be omitted if the primary key constraint is not enforced.

Note that intentionally I did not add end_date column as I found it not usable in my specific project's case.

What this statement does:

  • Create a stored procedure in repo schema with one parameter (table_name varchar(256)).
  • Dynamically obtains the table name you want to merge into from table_name parameter and performs UPDATE statement on is_active flag column (set to 0) on matching rows between dbimports table and its mirror table in repo schema.
  • Dynamically defines join conditions for the UPDATE statement.
  • Makes a INSERT * INTO from dbimports to repo schema with is_active flag set to 1 and current date (GETDATE()).

I know this is not a hard task if you build it in a normal SQL environment, but as you may know, Redshift's SQL is "highly modified". This means that most of normal SQL's features are cut out like variables and triggers, which makes this task pretty hard to figure out (at least for me it was).

Hope this helps anybody. Have a good one.

1

1 Answers

0
votes
    CREATE OR REPLACE PROCEDURE repo.inserter(table_name IN varchar(256))
AS $$
DECLARE 
    temp_table_name varchar := table_name || '_temp';
    one_liner varchar := temp_table_name || '_one_liner';
    query varchar;
    query_builder varchar;
    rec record;
BEGIN
EXECUTE 'DROP TABLE IF EXISTS ' || temp_table_name;
EXECUTE 'DROP TABLE IF EXISTS ' || one_liner;
EXECUTE 'CREATE TEMP TABLE ' || temp_table_name || ' (column_name varchar(256),ordinal_position int,is_nullable varchar(3));';
EXECUTE 'INSERT INTO ' || temp_table_name || '(column_name,ordinal_position,is_nullable) SELECT column_name,ordinal_position,is_nullable FROM svv_columns WHERE table_schema = ''repo'' AND table_name = ''' || table_name || ''';';
EXECUTE 'CREATE TEMP TABLE ' || one_liner || ' (columns_string varchar(1024));';
EXECUTE 'INSERT INTO ' || one_liner || '  SELECT LISTAGG(column_name,'','') WITHIN GROUP (order by ordinal_position) FROM ' || temp_table_name || ' where column_name not in ( ''scd_key'',''inserted_date'',''is_active'');';
query := 'SELECT column_name FROM ' || temp_table_name || ' WHERE is_nullable = ''NO'' and column_name not in(''scd_key'',''inserted_date'') order by ordinal_position;';
query_builder := 'UPDATE repo.' || table_name || ' SET is_active = 0 FROM dbimports.' || table_name || ' b WHERE ';
FOR rec IN EXECUTE query LOOP
query_builder := query_builder || 'repo.' || table_name || '.' || rec.column_name || ' = b.' || rec.column_name || ' AND ';
END LOOP;
query_builder := RTRIM(query_builder,'AND ') || ';';
EXECUTE query_builder;
query := 'SELECT columns_string FROM ' || one_liner || ';';
FOR rec IN EXECUTE query LOOP
query_builder := 'INSERT INTO repo.' || table_name || '(' ||  rec.columns_string || ',is_active,inserted_date) (SELECT ' || rec.columns_string || ',1 as is_active,GETDATE() as inserted_date FROM dbimports.' || table_name || ');';
END LOOP;
EXECUTE query_builder;

END;
$$  LANGUAGE plpgsql;