from airflow.providers.snowflake.hooks.snowflake import SnowflakeHook from jinja2 import Template from lib.config import SF_SCHEMA QUERY_TEMPLATE = Template(""" INSERT INTO ROYALTY_ACCOUNTING.{{schema}}.STMT_DB_SALES_DISTRO_STAGING ( UNIQUE_DETAIL_ID, BATCH_ID, STATEMENT_ID, START_DATE, TRANSACTION_DATE, STORE_ID, SUBDISTRIBUTOR, SALE_CURRENCY_CODE, SALE_CURRENCY_CODE_ID, ACTIVITY_RATE, COUNTRY_ID, CONFIGURATION, TRANSACTION_TYPE, TRANSACTION_TYPE_ID, TRANSACTION_SUBTYPE_ID, LABEL_ID, UPC, CD, TRACK_ID, ISRC, TRACK_NAME, VIDEO_ID, QUANTITY, UNIT_PRICE_USD, TOTAL_USD, WITHHOLDING_TAX_USD, RETAIL_PRICE_USD, ORIGINAL_PRICE_USD, DISCOUNT_USD, PHYS_PPD_USD, IS_EXCLUDED_FROM_SAP, ACTUAL_STATEMENT_NUMBER ) SELECT i.UNIQUE_DETAIL_ID, i.BATCH_ID, i.STATEMENT_ID, i.START_DATE, i.DATE, i.CUSTOMER_MASTER_MASTER_ID, i.SUBDISTRIBUTOR, i.ORIGINAL_CURRENCY_ISO, cm.ISO_NUMBER, i.ACTIVITY_RATE, i.COUNTRY_ID, i.CONFIGURATION, i.TRANS_TYPE, dtt.TRANSACTIONTYPEID, i.TRANS_SUBTYPE, i.VENDOR_ID, i.UPC, i.CD, i.TRACK_ID, i.ISRC, i.TRACK_NAME, i.VIDEO_ID, i.QTY, i.UNIT_PRICE, i.TOTAL, i.WHT, i.RETAIL_PRICE, i.ORIGINAL_PRICE, i.DISCOUNT, i.PHYS_PPD, i.SAP_EXCLUDE, i.ACTUAL_STATEMENT_NO FROM ROYALTY_ACCOUNTING.{{schema}}.STMT_DB_SALES_DISTRO_INGEST i LEFT JOIN FACTS.PROD.DIM_TRANSACTIONTYPE dtt ON dtt.TRANSACTIONTYPEABBR = i.TRANS_TYPE LEFT JOIN ROYALTY_ACCOUNTING.PROD.CURRENCY_MAP cm ON cm.ISO_CODE = i.ORIGINAL_CURRENCY_ISO """) def copy_sales(**kwargs): query = QUERY_TEMPLATE.render(schema=SF_SCHEMA) sf_hook = SnowflakeHook() sf_hook.run(query, autocommit=True)