In [1]:
import boto3
import pandas as pd
from io import StringIO, BytesIO
from datetime import datetime, timedelta

In [2]:
import config
import os

key_df = pd.read_csv('new_user_credentials.csv')
os.environ['AWS_ACCESS_KEY_ID'] = key_df['Access key ID'][0]
os.environ['AWS_SECRET_ACCESS_KEY'] = key_df['Secret access key'][0]

In [3]:
# Adapter Layer

def read_csv_to_df(bucket, key, decoding='utf-8', sep=','):
    csv_obj = bucket.Object(key=key).get().get('Body').read().decode(decoding)
    data = StringIO(csv_obj)
    df = pd.read_csv(data, delimiter=',')
    return df

def write_df_to_s3(bucket_target, df, key):
    out_buffer = BytesIO()
    df.to_parquet(out_buffer, index=False)
    bucket_target.put_object(Body=out_buffer.getvalue(), Key=key)
    return True

def list_files_in_prefix(bucket, prefix):
    files = [obj.key for obj in bucket.objects.filter(Prefix=prefix)]
    return files


#     objects = [obj for obj in bucket.objects.all()
#               if datetime.strptime(obj.key.split('/')[0], '%Y-%m-%d').date() >= min_date]
#     return objects

In [8]:
# Application Layer

def extract(bucket, date_list):
    files = [key for date in date_list for key in list_files_in_prefix(bucket, date)]
    df = pd.concat([read_csv_to_df(bucket, obj) for obj in files], ignore_index=True)
    return df

def transform_report1(df, columns, arg_date):
    
    df = df.loc[:, columns]
    df.dropna(inplace=True)
    df['opening_price'] = df.sort_values(by=['Time']).groupby(['ISIN', 'Date'])['StartPrice'].transform('first')
    df['closing_price'] = df.sort_values(by=['Time']).groupby(['ISIN', 'Date'])['StartPrice'].transform('last')
    df = df.groupby(['ISIN', 'Date'], as_index=False).agg(opening_price_eur=('opening_price', 'min'),
                                                             closing_price_eur=('closing_price', 'min'),
                                                             minimum_price_eur=('MinPrice', 'min'),
                                                             maximum_price_eur=('MaxPrice', 'max'),
                                                             daily_traded_volume=('TradedVolume', 'sum'))
    df['prev_closing_price'] = df.sort_values(by=['Date']).groupby(['ISIN'])['closing_price_eur'].shift(1)
    df['change_prev_closing_%'] = (df['closing_price_eur'] - df['prev_closing_price']) / df['prev_closing_price'] * 100
    df.drop(columns=['prev_closing_price'], inplace=True)
    df = df.round(decimals=2)
    df = df[df.Date >= arg_date]
    return df

def load(bucket, df, trg_key, trg_format):
    key =  trg_key + datetime.today().strftime("%Y%m%d_%H%M%S") + '.parquet' + trg_format
    print(bucket)
    print(key)
    write_df_to_s3(bucket, df, key)
    return True

def etl_report1(src_bucket, trg_bucket, date_list, columns, arg_date, trg_key, trg_format):
    df = extract(src_bucket, date_list)
    df = transform_report1(df, columns, arg_date)
    load(trg_bucket, df, trg_key, trg_format)
    return True

In [9]:
# Application Layer - not core
def return_date_list(bucket, arg_date, src_format):
    min_date = datetime.strptime(arg_date, '%Y-%m-%d').date() - timedelta(days=1)
    today = datetime.today().date()
    return_date_list = [(min_date + timedelta(days=x)).strftime(src_format) for x in range(0, (today-min_date).days + 1)]
    return return_date_list

In [10]:
# main function entrypoint
def main():
    # Parameters/Configurations
    # Later read config
    arg_date = '2021-08-04'
    src_format = '%Y-%m-%d'
    src_bucket = 'deutsche-boerse-xetra-pds'
    trg_bucket = 'xetra-1234'
    columns = ['ISIN', 'Date', 'Time', 'StartPrice', 'MaxPrice', 'MinPrice', 'EndPrice', 'TradedVolume']
    key = 'xetra_daily_report_' + datetime.today().strftime("%Y%m%d_%H%M%S") + '.parquet'
    trg_key = 'xetra_daily_report_'
    trg_format = '.parquet'
    
    # Init
    s3 = boto3.resource('s3')
    bucket_src = s3.Bucket(src_bucket)
    bucket_trg = s3.Bucket(trg_bucket)
    
    # run application
    objects = return_date_list(bucket_src, arg_date, src_format)
    etl_report1(bucket_src, bucket_trg, objects, columns, arg_date, trg_key, trg_format)

In [None]:
# run 
main()

In [None]:
trg_bucket = 'xetra-1234'
s3 = boto3.resource('s3')
bucket_trg = s3.Bucket(trg_bucket)
for obj in bucket_trg.objects.all():
    print(obj.key)

In [41]:
arg_date_dt = datetime.strptime(arg_date, '%Y-%m-%d').date() - timedelta(days=1)

In [42]:
arg_date_dt

datetime.date(2021, 8, 3)

In [43]:
s3 = boto3.resource('s3')
bucket = s3.Bucket('deutsche-boerse-xetra-pds')

In [44]:
objects = [obj for obj in bucket.objects.all()
          if datetime.strptime(obj.key.split('/')[0], '%Y-%m-%d').date() >= arg_date_dt]

In [45]:
def csv_to_df(filename):
    csv_obj = bucket.Object(key=filename).get().get('Body').read().decode('utf-8')
    data = StringIO(csv_obj)
    df = pd.read_csv(data, delimiter=',')
    return df

df_all = pd.concat([csv_to_df(obj.key) for obj in objects], ignore_index=True)

In [46]:
columns = ['ISIN', 'Date', 'Time', 'StartPrice', 'MaxPrice', 'MinPrice','EndPrice', 'TradedVolume']
df_all = df_all.loc[:, columns]

In [47]:
df_all.dropna(inplace=True)

## Get opening price per ISIN and day

In [48]:
df_all['opening_price'] = df_all.sort_values(by=['Time']).groupby(['ISIN', 'Date'])['StartPrice'].transform('first')

## Get closing price per ISIN and day

In [49]:
df_all['closing_price'] = df_all.sort_values(by=['Time']).groupby(['ISIN', 'Date'])['StartPrice'].transform('last')

## Aggregations

In [50]:
df_all = df_all.groupby(['ISIN', 'Date'], as_index=False).agg(opening_price_eur=('opening_price', 'min'),
                                                             closing_price_eur=('closing_price', 'min'),
                                                             minimum_price_eur=('MinPrice', 'min'),
                                                             maximum_price_eur=('MaxPrice', 'max'),
                                                             daily_traded_volume=('TradedVolume', 'sum'))

## Percent Change Prev Closing

In [51]:
df_all['prev_closing_price'] = df_all.sort_values(by=['Date']).groupby(['ISIN'])['closing_price_eur'].shift(1)
df_all['change_prev_closing_%'] = (df_all['closing_price_eur'] - df_all['prev_closing_price']) / df_all['prev_closing_price'] * 100
df_all.drop(columns=['prev_closing_price'], inplace=True)

In [54]:
df_all = df_all.round(decimals=2)

In [55]:
df_all = df_all[df_all.Date >= arg_date]

## Write to S3

In [56]:
key = 'xetra_daily_report_' + datetime.today().strftime("%Y%m%d_%H%M%S") + '.parquet'

In [57]:
out_buffer = BytesIO()
df_all.to_parquet(out_buffer, index=False)
bucket_target = s3.Bucket('xetra-12345')
bucket_target.put_object(Body=out_buffer.getvalue(), Key=key)

s3.Object(bucket_name='xetra-12345', key='xetra_daily_report_20210805_065759.parquet')

# Reading the uploaded file

In [58]:
for obj in bucket_target.objects.all():
    print(obj.key)

xetra_daily_report_20210805_065412.parquet
xetra_daily_report_20210805_065759.parquet


In [61]:
prq_obj = bucket_target.Object(key='xetra_daily_report_20210805_065759.parquet').get().get('Body').read()
data = BytesIO(prq_obj)
df_report = pd.read_parquet(data)

In [62]:
df_report

Unnamed: 0,ISIN,Date,opening_price_eur,closing_price_eur,minimum_price_eur,maximum_price_eur,daily_traded_volume,change_prev_closing_%
0,AT00000FACC2,2021-08-04,8.69,8.48,8.48,8.75,721,-1.40
1,AT0000606306,2021-08-04,19.94,19.91,19.72,19.94,1269,0.81
2,AT0000609607,2021-08-04,16.68,16.60,16.60,16.68,786,0.36
3,AT0000644505,2021-08-04,110.80,110.80,110.80,110.80,4,2.03
4,AT0000652011,2021-08-04,33.50,33.50,33.50,33.50,0,1.21
...,...,...,...,...,...,...,...,...
3029,XS2265368097,2021-08-04,15.27,15.26,15.26,15.38,2700,0.00
3030,XS2265369574,2021-08-04,21.72,21.51,21.51,21.74,6,-0.39
3031,XS2265369731,2021-08-04,8.86,8.72,8.72,8.86,106,-1.27
3032,XS2265370234,2021-08-04,22.46,22.41,22.41,22.47,300,0.48
