<a id="CELL1"></a>
## CELL 1 


In [1]:
"""CELL 1
builds and returns a database session
local assumes a psql instance in a local docker container
only postgres database is supported for configuration_application at this time
"""
"""
gets env-based configuration secret
returns a session to the configuration db
for dev env it pre-populates the database with helper and seed data
"""
from core.helpers.session_helper import SessionHelper
session = SessionHelper().session

2019-08-06 11:57:47,221 - core.helpers.session_helper.SessionHelper - INFO - Creating session for dev environment...
2019-08-06 11:57:47,246 - core.helpers.configuration_mocker.ConfigurationMocker - DEBUG - Generating administrator mocks.
2019-08-06 11:57:47,284 - core.helpers.configuration_mocker.ConfigurationMocker - DEBUG - Done generating administrator mocks.
2019-08-06 11:57:47,285 - core.helpers.configuration_mocker.ConfigurationMocker - DEBUG - Generating pharmaceutical company mocks.
2019-08-06 11:57:47,289 - core.helpers.configuration_mocker.ConfigurationMocker - DEBUG - Done generating pharmaceutical company mocks.
2019-08-06 11:57:47,290 - core.helpers.configuration_mocker.ConfigurationMocker - DEBUG - Generating brand mocks.
2019-08-06 11:57:47,297 - core.helpers.configuration_mocker.ConfigurationMocker - DEBUG - Done generating brand mocks.
2019-08-06 11:57:47,298 - core.helpers.configuration_mocker.ConfigurationMocker - DEBUG - Generating segment mocks.
2019-08-06 11:57:4

## CONFIGURATION - PLEASE TOUCH
### <font color=pink>This cell will be off in production as configurations will come from the configuration postgres DB</color>

In [3]:
"""
************ CONFIGURATION - PLEASE TOUCH **************
Pipeline Builder configuration: creates configurations from variables specified here!!
This cell will be off in production as configurations will come from the configuration postgres DB.
"""
"""
PIPELINE STATE:

raw-->ingest-->master-->enhance-->enrich-->metrics-->dimensional

"""
# config vars: this dataset
config_pharma = "sun" # the pharmaceutical company which owns {brand}
config_brand = "ilumya" # the brand this pipeline operates on
config_state = "master" # the state this transform runs in
config_name = "master_referral_source" # the name of this transform, which is the name of this notebook without .ipynb

# input vars: dataset to fetch. 
# Recall that a contract published to S3 has a key format branch/pharma/brand/state/name
input_branch = "sun-extract-validation"
# None
# if None, input_branch is automagically set to your working branch
input_pharma = "sun"
input_brand = "ilumya"
input_state = "ingest"
input_name = "symphony_health_association_ingest_column_mapping"

#This contract defines the base of the output structure of data into S3.
#
#contract structure in s3: 
#s3:// {ENV} / {BRANCH} / {PARENT} / {CHILD} / {STATE} / {name of input}
#
#ENV - environment Must be one of development, uat, production.
#Prefixed with integrichain- due to global unique reqirement
#BRANCH - the software branch for development this will be the working pull request (eg pr-225)
#in uat this will be edge, in production this will be master
#PARENT - The top level source identifier
#this is generally the customer (and it is aliased as such) but can be IntegriChain for internal sources,
#or another aggregator for future-proofing
#CHILD - The sub level source identifier, generally the brand (and is aliased as such)
#STATE - One of: raw, ingest, master, enhance, enrich, metrics


In [4]:
import logging
logging.getLogger().setLevel(logging.DEBUG)
log = logging.getLogger()

### <font color=orange>SETUP - DON'T TOUCH </font>
Populating config mocker based on config parameters...

In [5]:
"""
************ SETUP - DON'T TOUCH **************
Populating config mocker based on config parameters...
"""
import core.helpers.pipeline_builder as builder

ids = builder.build(config_pharma, config_brand, config_state, config_name, session)
"""
RETURNS: A list of 2 items: [transformation_id, run_id] where transformation_id corresponds
to the configuration created/found for {transformation} and run_id is a randomly generated 6 digit
number (to avoid publishing to the same place with the same dataset)
"""
transform_id = ids[0]
run_id = ids[1]

2019-08-06 11:59:49,196 - core.logging - DEBUG - Adding/getting mocks for specified configurations...
2019-08-06 11:59:49,221 - core.logging - DEBUG - Done. Creating mock run event and committing results to configuration mocker.


In [6]:
# debug only
# e.g.
# 6
# 136126
log.debug('Transform Id:{} Run Id:{}'.format(transform_id,run_id))

2019-08-06 12:00:02,075 - root - DEBUG - Transform Id:6 Run Id:300953


### <font color=orange>SETUP - DON'T TOUCH </font>
This section imports data from the configuration database
and should not need to be altered or otherwise messed with. 


In [7]:
"""************ SETUP - DON'T TOUCH **************
This section imports data from the configuration database
and should not need to be altered or otherwise messed with. 
~~These are not the droids you are looking for~~
"""
from core.constants import BRANCH_NAME, ENV_BUCKET, BATCH_JOB_QUEUE
from core.helpers.session_helper import SessionHelper
from core.models.configuration import Transformation
from dataclasses import dataclass
from core.dataset_contract import DatasetContract


db_transform = session.query(Transformation).filter(Transformation.id == transform_id).one()

@dataclass
class DbTransform:
    id: int = db_transform.id ## the instance id of the transform in the config app
    name: str = db_transform.transformation_template.name ## the transform name in the config app
    state: str = db_transform.pipeline_state.pipeline_state_type.name ## the pipeline state, one of raw, ingest, master, enhance, enrich, metrics, dimensional
    branch:str = BRANCH_NAME ## the git branch for this execution 
    brand: str = db_transform.pipeline_state.pipeline.brand.name ## the pharma brand name
    pharmaceutical_company: str = db_transform.pipeline_state.pipeline.brand.pharmaceutical_company.name # the pharma company name
    publish_contract: DatasetContract = DatasetContract(branch=BRANCH_NAME,
                            state=db_transform.pipeline_state.pipeline_state_type.name,
                            parent=db_transform.pipeline_state.pipeline.brand.pharmaceutical_company.name,
                            child=db_transform.pipeline_state.pipeline.brand.name,
                            dataset=db_transform.transformation_template.name)


In [8]:
#debug only
# e.g:
#DC-578_PatientStatus
#ichain-dev
#dev-core
log.debug('Branch name:{} Env Bucket:{} Batch Job Queue:{}'.format(BRANCH_NAME,ENV_BUCKET,BATCH_JOB_QUEUE))

2019-08-06 12:00:32,962 - root - DEBUG - Branch name:DC-580_ReferralSource Env Bucket:ichain-dev Batch Job Queue:dev-core


***
# CORE Cartridge Notebook::[master_referral_source]
![CORE Logo](assets/coreLogo.png) 

---
## Keep in Mind
Good Transforms Are...
- **singular in purpose:** good transforms do one and only one thing, and handle all known cases for that thing. 
- **repeatable:** transforms should be written in a way that they can be run against the same dataset an infinate number of times and get the same result every time. 
- **easy to read:** 99 times out of 100, readable, clear code that runs a little slower is more valuable than a mess that runs quickly. 
- **No 'magic numbers':** if a variable or function is not instantly obvious as to what it is or does, without context, maybe consider renaming it.

## Workflow - how to use this notebook to make science
#### Data Science
1. **Document your transform.** Fill out the _description_ cell below describing what it is this transform does; this will appear in the configuration application where Ops will create, configure and update pipelines. 
1. **Define your config object.** Fill out the _configuration_ cell below the commented-out guide to define the variables you want ops to set in the configuration application (these will populate here for every pipeline). 
2. **Build your transformation logic.** Use the transformation cell to do that magic that you do. 
![caution](assets/cautionTape.png)

## CONFIGURATION - VARIABLES - PLEASE TOUCH

# TRANSFORM

In [37]:
""" 
CONFIGURATION ********* VARIABLES - PLEASE TOUCH ********* 
This section defines what you expect to get from the configuration application 
in a single "transform" object. Define the vars you need here, and comment inline to the right of them 
for all-in-one documentation. 
Engineering will build a production "transform" object for every pipeline that matches what you define here.

@@@ FORMAT OF THE DATA CLASS IS: @@@ 

<variable_name>: <data_type> #<comment explaining what the value is to future us>
e.g.
class Transform(DbTransform):
    some_ratio: float
    site_name: str

~~These ARE the droids you are looking for~~
"""
"""
imports
"""
import pandas as pd
from core.logging import get_logger
 
class Transform(DbTransform):
    '''
    YOUR properties go here!!
    Variable properties should be assigned to the exact name of
    the transformation as it appears in the Jupyter notebook filename.
    ''' 

    col_referral_source: str           
    
    
    def master_referral_source(self,df):
        try:        
            go = False # assume things are not working YET.
           
            dffail = pd.DataFrame() # initialize df for fails
            
            # master data golden referral
            #working draft] Gold Domain of referral source
            logger.info('try:')
            referral_dict = {}
            referral_dict[1]='DIRECT'
            referral_dict[2]='HUB'
            referral_dict[3]='PHARM'
            
            # store the golden values in a list
            referral_list = list(referral_dict.values())           
            
            # log metadata 
            
            logger.info('Gold Domain List:{}'.format(referral_list))  
            
            # df in
            dfShape = df.shape
            logger.info('df in  shape: {} {}'.format(dfShape[0],dfShape[1])) 
            logger.info('df in {}'.format(df.head()))   
            
            # am I expecting certain column names? YES 
            referralColNameExpected = transform.col_referral_source
            
            logger.info('expecting column name  referral as:{}'.format(referralColNameExpected))
            columnNamesArr = df.columns.values.tolist()
            logger.info('df column names:{}'.format(columnNamesArr))
            
            if referralColNameExpected in columnNamesArr:
 
                # apply Upper Case to col values of interest
                df[referralColNameExpected]= df[referralColNameExpected].apply(lambda x: x.upper() if x is not None else x)   
                # apply strip  to col values of interest
                df[referralColNameExpected]= df[referralColNameExpected].apply(lambda x: x.strip() if x is not None else x)
                # what fails
                dffail = df[~df[referralColNameExpected].isin(referral_list)]
                # apply master selection for the column of interest
                # what passes
                df = df[df[referralColNameExpected].isin(referral_list)]
                
                # meta data log for what comes out of the function pass and fail df
                dfOutSize = df.size
                dfOutShape = df.shape
                dffailSize = dffail.size
                dffailShape = dffail.shape
                logger.info('df in   shape: {} {}'.format(dfShape[0],dfShape[1]))                 
                logger.info('df pass shape: {} {}'.format(dfOutShape[0],dfOutShape[1]))
                logger.info('df fail shape: {} {}'.format(dffailShape[0],dffailShape[1]))
                logger.info('df pass {}'.format(df.head()))
                logger.info('df fail {}'.format(dffail.head()))   
                # end meta data
                go = True
            else:
                go = False # something did not work
                logger.exception('expecting column name for referral_source if/else exception raise')
                raise Exception("master_referral_source try if/else referralColNameExpected in columnNamesArr")              
        except Exception as e:
            go = False # something did not work
            logger.exception("exception:".format(e))
            raise Exception(str(e))
        else:
            pass
        finally:
            pass
        return df.copy(),dffail.copy(),go
                
transform = Transform()
logger = get_logger(f"core.transforms.{transform.state}.{transform.name}")

### *Please place your value assignments for development below*
### <font color=pink>This cell will be turned off in production, Engineering will set to pull from the configuration</color>

In [38]:
## Please place your value assignments for development here!!
## This cell will be turned off in production and Engineering will set to pull from the configuration application instead
## For the last example, this could look like...
## transform.some_ratio = 0.6
## transform.site_name = "WALGREENS"

transform.col_referral_source = 'ref_source'


### Description
What does this transformation do? be specific.

![what does your transform do](assets/what.gif)

## Planned
Need a transform that is able to map multiple distinct instances of a referral source to a cleansed referral source data model.

Definition of Done:
- Collect all unique raw referral source instances
- Auto-map as many raw referral source instances to a defined cleansed data model.
- Process for identifying and manually mapping where auto-map fails.
= Do not publish un-mapped instances. Drop them, give us the ability to triage and map to gold.


### FETCH DATA - TOUCH, BUT CAREFULLY
### <font color=pink>This cell will be turned off in production, as the input_contract will be handled by the pipeline</color>

In [11]:
"""
************ FETCH DATA - TOUCH, BUT CAREFULLY **************
This cell will be turned off in production, as the input_contract will be handled by the pipeline.
"""
logger.info("FETCH DATA CELL - TOUCH - This cell will be turned off in production, as the input_contract will be handled by the pipeline. ")

# for testing / development only !!! picking run id based on core pipeline datasets available
run_id = 3

if not input_branch:
    input_branch = BRANCH_NAME
input_contract = DatasetContract(branch=input_branch,
                                 state=input_state, 
                                 parent=input_pharma, 
                                 child=input_brand, 
                                 dataset=input_name)
run_filter = []
run_filter.append(dict(partition="__metadata_run_id", comparison="==", values=[run_id]))
# IF YOU HAVE PUBLISHED DATA MULTIPLE TIMES, uncomment the above line and change the int to the run_id to fetch.
# Otherwise, you will have duplicate values in your fetched dataset!

# bypass/comment out when unit testing individual parquet files
df = input_contract.fetch(filters=run_filter)



2019-08-06 12:32:39,706 - core.transforms.master.master_referral_source - INFO - FETCH DATA CELL - TOUCH - This cell will be turned off in production, as the input_contract will be handled by the pipeline. 
2019-08-06 12:32:39,710 - core.dataset_contract.DatasetContract - INFO - Fetching dataframe from s3 location s3://ichain-dev/sun-extract-validation/sun/ilumya/ingest/symphony_health_association_ingest_column_mapping.
2019-08-06 12:32:39,937 - urllib3.util.retry - DEBUG - Converted retries value: False -> Retry(total=False, connect=None, read=None, redirect=0, status=None)
2019-08-06 12:32:39,939 - urllib3.connectionpool - DEBUG - Starting new HTTPS connection (1): ichain-dev.s3.amazonaws.com:443
2019-08-06 12:32:40,413 - urllib3.connectionpool - DEBUG - https://ichain-dev.s3.amazonaws.com:443 "GET /?prefix=sun-extract-validation%2Fsun%2Filumya%2Fingest%2Fsymphony_health_association_ingest_column_mapping&encoding-type=url HTTP/1.1" 200 None
2019-08-06 12:32:40,436 - urllib3.util.retr

## *<font color=grey>unit test development only*</font>
*<font color=grey>The next **5** cells will be deleted in production.* </font>

In [None]:
#import pyarrow.parquet as pq
#import s3fs

#def pandas_from_parquet_s3(file_path):  
#    s3 = s3fs.S3FileSystem()
#    df = (
#        pq
#        .ParquetDataset(file_path, filesystem=s3)
#        .read_pandas()
#        .to_pandas()
#    )    
#    return df

In [None]:
# unit test/development 
# isolate on individual parquet files
#TEST 1
#df = pandas_from_parquet_s3('ichain-dev/sun-extract-validation/sun/ilumya/ingest/symphony_health_association_ingest_column_mapping/__metadata_run_id=3/90ca3aa7b0bb4246a281591b013ff54e.parquet')
# TEST 2
#df = pandas_from_parquet_s3('ichain-dev/sun-extract-validation/sun/ilumya/ingest/symphony_health_association_ingest_column_mapping/__metadata_run_id=3/d7ad974cef284e19aa7b5ac410220b96.parquet')
# TEST 3
#df = pandas_from_parquet_s3('ichain-dev/sun-extract-validation/sun/ilumya/ingest/symphony_health_association_ingest_column_mapping/__metadata_run_id=3/1a6ffd3598d442e38fbba66ea85a55a2.parquet')
# TEST 4
#df = pandas_from_parquet_s3('ichain-dev/sun-extract-validation/sun/ilumya/ingest/symphony_health_association_ingest_column_mapping/__metadata_run_id=3/5c00059d9fc04b0e8bc4ce764c50f3fb.parquet')
# TEST 5
# df = pandas_from_parquet_s3('ichain-dev/sun-extract-validation/sun/ilumya/ingest/symphony_health_association_ingest_column_mapping/__metadata_run_id=3/6eceb7ce59bd4dec8720316b4209b0e3.parquet')
# THEN ALL TEST use 
# then use the FETCH DATA - TOUCH, BUT CAREFULLY CELL

In [None]:
# unit test/development only
# before shot unit testing only
# dfSize = df.size
# dfShape = df.shape
# print('shape: {} {}'.format(dfShape[0],dfShape[1])) 

In [None]:
# unit test/development only
# needed to see the col(s) of interest
#pd.set_option('display.max_columns', 50)

In [30]:
## unit test/development only
#df.head()
#df[transform.col_referral_source].value_counts(dropna=False)

DIRECT    15186
HUB        8689
PHARM       572
NaN          10
Name: ref_source, dtype: int64

# <font color=red>**CALL**</font> THE TRANSFORM

In [39]:
### Use the variables above to execute your transformation.
### the final output needs to be a variable named final_dataframe
logger.info("CALL THE TRANSFORM - execute your transformation - the final output needs to be a variable named final_dataframe")

final_dataframe, final_fail, go = transform.master_referral_source(df)

if go==True:
    logger.info("CALL THE TRANSFORM -  go no go = GO")
elif go==False:
    logger.info("CALL THE TRANSFORM -  go no go = NO go")
else:
    go=False
    logger.info("CALL THE TRANSFORM -  go no go = unknown make it NO go")
        

2019-08-06 13:19:14,983 - core.transforms.master.master_referral_source - INFO - CALL THE TRANSFORM - execute your transformation - the final output needs to be a variable named final_dataframe
2019-08-06 13:19:14,987 - core.transforms.master.master_referral_source - INFO - try:
2019-08-06 13:19:14,989 - core.transforms.master.master_referral_source - INFO - Gold Domain List:['DIRECT', 'HUB', 'PHARM']
2019-08-06 13:19:14,990 - core.transforms.master.master_referral_source - INFO - df in  shape: 24457 72
2019-08-06 13:19:15,026 - core.transforms.master.master_referral_source - INFO - df in          rec_date pharm_code   pharm_npi transtype pharm_transaction_id  \
0  20181024115959    ACCREDO  1346208949       COM   279133432018102401   
1  20181025115959    ACCREDO  1346208949       COM   278370982018102502   
2  20181029115959    ACCREDO  1346208949       COM   279181482018102903   
3  20181102115959    ACCREDO  1346208949       COM   267244982018110204   
4  20181106115959    ACCREDO 

### *<font color=grey>unittest python*</font>

In [33]:
import unittest

def ut_shape(final_dataframe,df):
    """
    assertion will change based on coding state
    """
    return final_dataframe.shape == df.shape

class TestNotebook(unittest.TestCase):
    
    def test_ut_shape(self):
        
        self.assertEqual(ut_shape(final_dataframe,df),True)
                
"""

"""
# for development only
unittest.main(argv=[''], verbosity= 2, exit=False)        
    

test_ut_shape (__main__.TestNotebook) ... FAIL

FAIL: test_ut_shape (__main__.TestNotebook)
----------------------------------------------------------------------
Traceback (most recent call last):
  File "<ipython-input-33-f771f8d9cda6>", line 13, in test_ut_shape
    self.assertEqual(ut_shape(final_dataframe,df),True)
AssertionError: False != True

----------------------------------------------------------------------
Ran 1 test in 0.008s

FAILED (failures=1)


<unittest.main.TestProgram at 0x7f9c6d6197f0>

In [35]:
# untit test/development only look at the fails
#final_fail.head()
#final_fail[transform.col_referral_source]

6      None
7      None
8      None
9      None
10     None
712    None
713    None
714    None
715    None
716    None
Name: ref_source, dtype: object

In [None]:
# untit test/development only look at the pass(es)
#final_dataframe.head()
#final_dataframe[transform.col_referral_source]

# **publish**
### Writing to S3
Invoke the `publish()` command to write to a given contract. Some things to know:
- To invoke publish a contract must be at the grain of dataset. This is because file names will be set by the dataframe=\>parquet conversion. 
- publish only accepts a pandas dataframe.
- publish does not allow for timedelta data types at this time (this is missing functionality in pyarrow).
- publish handles partitioning the data as per contract, creating file paths, and creating the binary parquet files in S3, as well as the needed metadata. <br>
**- by default, all datasets include a single partition, \_\_metadata\_run\_id, the RunEvent ID of an executed pipeline**

In [36]:
## that's it - just provide the final dataframe to the var final_dataframe and we take it from there
if go==True:
    logger.info("PUBLISH - that's it - its a GO - just provide the final dataframe to the var final_dataframe and we take it from there")
    transform.publish_contract.publish(final_dataframe, run_id, session)
elif go==False:
    logger.info("PUBLISH -  go no go = NO go -  so DONT publish")
else:
    go=False
    logger.info("PUBLISH -  go no go = unknown make it NO go - so DONT publish")    
session.close()

2019-08-06 13:02:39,203 - core.transforms.master.master_referral_source - INFO - PUBLISH - that's it - its a GO - just provide the final dataframe to the var final_dataframe and we take it from there
2019-08-06 13:02:39,205 - core.dataset_contract.DatasetContract - INFO - Publishing dataframe to s3 location s3://ichain-dev/dc-580_referralsource/sun/ilumya/master/master_referral_source with run ID 3.
2019-08-06 13:02:39,214 - core.dataset_contract.DatasetContract - DEBUG - Publishing dataframe to Redshift Spectrum database ichain_core to schema.table                 data_core.sun_ilumya_master_referral_source...
2019-08-06 13:02:39,217 - s3parq.publish_parq - DEBUG - Found redshift parameters. Checking validity of params...
2019-08-06 13:02:39,222 - s3parq.publish_parq - DEBUG - Checking redshift params are correctly formatted
2019-08-06 13:02:39,227 - s3parq.publish_parq - DEBUG - Done checking redshift params
2019-08-06 13:02:39,232 - s3parq.publish_parq - DEBUG - Redshift parameters 

  """)


ProgrammingError: (psycopg2.ProgrammingError) permission denied for database ichain_core

[SQL: CREATE EXTERNAL SCHEMA IF NOT EXISTS data_core                 FROM DATA CATALOG                 database 'ichain_core'                 iam_role 'arn:aws:iam::265991248033:role/mySpectrumRole';]
(Background on this error at: http://sqlalche.me/e/f405)

***