diff --git a/.gitignore b/.gitignore index 16f2dc5..a9c44ae 100644 --- a/.gitignore +++ b/.gitignore @@ -1 +1,2 @@ -*.csv \ No newline at end of file +*.csv +__pycache__ \ No newline at end of file diff --git a/README.md b/README.md index 00c3daf..18ed5fa 100644 --- a/README.md +++ b/README.md @@ -7,4 +7,11 @@ "orders" : "SELECT * FROM orders;", "payments" : "SELECT * FROM payments;", "productlines" : "SELECT * FROM productlines;", - "products" : "SELECT * FROM products;" \ No newline at end of file + "products" : "SELECT * FROM products;" + + + { + "transactions" : "SELECT * FROM transactions;", + "bookings" : "SELECT * FROM bookings;", + "service_types" : "SELECT * FROM service_types;" +} \ No newline at end of file diff --git a/config/CLOUD-STORAGE-ADMIN-CREDENTIAL.json b/config/CLOUD-STORAGE-ADMIN-CREDENTIAL.json new file mode 100644 index 0000000..909b7fd --- /dev/null +++ b/config/CLOUD-STORAGE-ADMIN-CREDENTIAL.json @@ -0,0 +1,12 @@ +{ + "type": "service_account", + "project_id": "data-management-359208", + "private_key_id": "06a6c2d224dff316c66498d62515b35ebdbe946f", + "private_key": "-----BEGIN PRIVATE KEY-----\nMIIEvgIBADANBgkqhkiG9w0BAQEFAASCBKgwggSkAgEAAoIBAQDbnWVbVtvG5Cbt\n8p3eriMFcQeFqFS1xp4srM09pFqmnvYMOJw8aJuK7RoRNlZQWKTgz8MabMGYaakL\nJJeiI6bUZO8EDAkui/TyiAseqJjOPBTgavHklyUKFSy/hdzLqRseEmi8SD1VU/Ek\nlzMARsnXw0Zrd+eptsR2MuKWQX4Qjj0FH1FDW/3GP2jgzS+HRktPzT3QIHlhyAOV\nxaCYoRY1IX/U8eNf0ZIu8SZzErLEjZRdhjMRuQVsvH4YZ11zbpI7UZZ/eLjT6zOB\nvbYJO6hXWFdTSej9FBg6e1ZJ82zVeMcFkaoNsfkD2CQqdv/3pXUbgtBLGCBuhNu3\nwQHyORIHAgMBAAECggEAaIm3oY7q9vXLgiCm/USu7vwqtHi4Of7ddC6dU+ZUMFQi\nkxavaCHzSGIsslzHIV/QvCKpoH58eOxyxxcYBtopo5iYHbkM9dcxNfGEOYfPlPwM\ng/bkRgecXfxOXKx/uYI5okrpCBbq+x8F/oDqigsoMUiG0Mk2wRZ61jjKmvN56q6o\nam59GBQoLOqh5uWWq/X0YqvsE8PaXqPvLwEZgeqTID6m2HAT6C4WSf+zuLe1vZ5O\nT19ue7ypWkgGD724ChQ/p/I9dy9P2nmtKDQ+4yGnHh8RU9WuXiSwRbDVJ6eV/C5J\nM27V1SoQmjAFYoQU+bq9tu0mHFPmiVELAVUPDW5oQQKBgQD4ZaBgLaRnGCBxfEBQ\n16B0Ks1NuxxSdyRk2aJ2OcDIXTLTcI1t5dJRgF836XqEoIHDrXaduVfhOqnbuJWv\nvcpbNhn3762hLCaih1bbz13GWIBqvVS8oJojJezC2r4tYgM3ydynx9C24YsBIhBl\nU2MBSufXUeDfI6ir+MLyNGuqjwKBgQDiVj111vqU7oABheeE1QewJMR4PdnKZWJ7\nGJHH1sBbaOeO0xYSFMg32fKvWLC2hEFJiNqWy90k7SmTcOPES2sV7nREk9MLO6Fy\nZyN4gw+5n3g+n66VAgTq2i99/vD95ubLH85GBfAt2KM07y3rgcNcEaP1QI9ue2lM\nHrVzNnQ9CQKBgQDOIvBTwKzljWUnKLjrHfafUQHtlvDrEsqWEvI68LSm0okiZQ5J\nfGbskf7zFIRDWjw2GlcMj0p5tEhP+j/mhzdOOHiWhEXwMgah7HTNl6o3tyxi6FpQ\n62re7lMsZYFbgjIvcwr2BeGUU1obB5zZqbjI0tPRobZfF2WbyaZmf9A1ywKBgQDg\nJ3q83sjSgJWjbIMKqZPwnak6UD8GVHxA3udZm9RrcyyI5YLhK1XTAnV3tQVl7Ptf\noTqix4nfTUW0sMPSHsMSOFNLq38Ci+7rhzu42UvUkRucIbbb+eD22ljYlokDXA9M\nMdauwKjKLtgLz6iRqbTZ1NqlRGgIig6RhYQ8czyRSQKBgDi+eLyxt+7zkyXMNhcO\nPEI4lOmjgSsHIKHByKlNx0EofmQYIu0zGwFWXd4MLApbfBy8zT2AqTwhv8XdhCn3\nsjwdNxpdHx1aJHFXkH1hndzzPz6Ud9oMMpTEcAgoYAoB6YV6OpLKQSGBUpbFzWdM\nYSuH1Vl4SdfL0p8/yS0UH7ed\n-----END PRIVATE KEY-----\n", + "client_email": "cloud-storage-admin-service-ac@data-management-359208.iam.gserviceaccount.com", + "client_id": "111008631797803896886", + "auth_uri": "https://accounts.google.com/o/oauth2/auth", + "token_uri": "https://oauth2.googleapis.com/token", + "auth_provider_x509_cert_url": "https://www.googleapis.com/oauth2/v1/certs", + "client_x509_cert_url": "https://www.googleapis.com/robot/v1/metadata/x509/cloud-storage-admin-service-ac%40data-management-359208.iam.gserviceaccount.com" +} diff --git a/python_el/extract/exchange_extract/EXCHANGE_CONFIG.json b/config/EXCHANGE_CONFIG.json similarity index 100% rename from python_el/extract/exchange_extract/EXCHANGE_CONFIG.json rename to config/EXCHANGE_CONFIG.json diff --git a/config/EXCHANGE_EXTRACT_SQL.json b/config/EXCHANGE_EXTRACT_SQL.json new file mode 100644 index 0000000..fc7a7e8 --- /dev/null +++ b/config/EXCHANGE_EXTRACT_SQL.json @@ -0,0 +1,10 @@ +{ + "customers" : "SELECT * FROM customers;", + "employees" : "SELECT * FROM employees;", + "offices" : "SELECT * FROM offices;", + "orderdetails" : "SELECT * FROM orderdetails;", + "orders" : "SELECT * FROM orders;", + "payments" : "SELECT * FROM payments;", + "productlines" : "SELECT * FROM productlines;", + "products" : "SELECT * FROM products;" +} \ No newline at end of file diff --git a/exchange/extract/exchange_extract_function.py b/exchange/extract/exchange_extract_function.py new file mode 100644 index 0000000..88b4601 --- /dev/null +++ b/exchange/extract/exchange_extract_function.py @@ -0,0 +1,45 @@ +import json +import csv +from operator import concat +import mysql.connector + + +def connect_mysql(config): + try: + mydb = mysql.connector.connect( + host = config["host"], + user = config["user"], + password = config["password"], + database = config["database"]) + return mydb + except Exception as e: + print(e) + print("[FUNCTION_ERROR]-read_json_file:" , name) + return False + + +def mysql_dict_sql_to_csv(SQL_DICT, path, mysqldb, mysqlcursor): + + for tablename,sql in SQL_DICT.items(): + + try: + mysqlcursor.execute(sql) + mysqlresult = mysqlcursor.fetchall() + mysqldbname = mysqldb.database + + filename = '/{}_{}.csv'.format(mysqldbname, tablename) + headers = [col[0] for col in mysqlcursor.description] # get headers + mysqlresult.insert(0, tuple(headers)) + fp = open(concat(path,filename), 'w', newline = '') + myFile = csv.writer(fp) + myFile.writerows(mysqlresult) + fp.close() + + except Exception as e: + print(e) + print("[FUNCTION_ERROR]-write_to_csv:" , filename) + return False + + finally: + print("[Download Complete]",filename) + diff --git a/exchange/load/exchange_load_function.py b/exchange/load/exchange_load_function.py new file mode 100644 index 0000000..09ff08f --- /dev/null +++ b/exchange/load/exchange_load_function.py @@ -0,0 +1,23 @@ +import os +from google.cloud import storage +import json + +def login_cloudstorage_credential(credential): + try: + storage_client = storage.Client.from_service_account_json(credential) + return storage_client + except Exception as e: + print(e) + print("[FUNCTION_ERROR]-login_cloudstorage_credential") + return False + +def upload_to_bucket(storage_client, blob_name, file_path): + try: + exchange_bucket = storage_client.get_bucket('exchange_bucket_raw') + blob = exchange_bucket.blob(blob_name) + blob.upload_from_filename(file_path) + return print("[Upload Complete]",file_path) + except Exception as e: + print(e) + print("[FUNCTION_ERROR]-upload_to_bucket",blob_name) + return False diff --git a/general_function.py b/general_function.py new file mode 100644 index 0000000..51c5d49 --- /dev/null +++ b/general_function.py @@ -0,0 +1,12 @@ +import json +from operator import concat + +def read_json_file(path, name): + try: + with open(concat(path,name)) as json_file: + output = json.load(json_file) + return output + except Exception as e: + print(e) + print("[FUNCTION_ERROR]-read_json_file:" , name) + return False diff --git a/python_el/extract/exchange_extract/EXCHANGE_EXTRACT_SQL.json b/python_el/extract/exchange_extract/EXCHANGE_EXTRACT_SQL.json deleted file mode 100644 index cea1bb2..0000000 --- a/python_el/extract/exchange_extract/EXCHANGE_EXTRACT_SQL.json +++ /dev/null @@ -1,5 +0,0 @@ -{ - "transactions" : "SELECT * FROM transactions;", - "bookings" : "SELECT * FROM bookings;", - "service_types" : "SELECT * FROM service_types;" -} \ No newline at end of file diff --git a/python_el/extract/exchange_extract/exchange_extract_main.py b/python_el/extract/exchange_extract/exchange_extract_main.py deleted file mode 100644 index 75f8689..0000000 --- a/python_el/extract/exchange_extract/exchange_extract_main.py +++ /dev/null @@ -1,91 +0,0 @@ -from operator import concat -import mysql.connector -import json -import os -import csv - -#Get all the Compulsory File -CURRENT_DIRECTORY_PATH = "\\python_el\\extract\\exchange_extract" -MYSQL_CONFIG_FILE = "EXCHANGE_CONFIG.json" -SQL_FILE = "EXCHANGE_EXTRACT_SQL.json" -CSV_STORAGE_PATH = "exchange_temp_extract_csv/" - -#Function -def change_directory_to(path): - try: - current_directory = os.getcwd() - target_dir = concat(current_directory,path) - return os.chdir(target_dir) - - except Exception as e: - print(e) - print("[FUNCTION_ERROR]-change_directory_to") - - - -def read_json_file(name): - try: - with open(name) as json_file: - output = json.load(json_file) - return output - except Exception as e: - print(e) - print("[FUNCTION_ERROR]-read_json_file:" , name) - - - -def connect_mysql(config): - mydb = mysql.connector.connect( - host = config["host"], - user = config["user"], - password = config["password"], - database = config["database"]) - return mydb - - - -def write_to_csv(path,dbname,tablename,cursor,data): - try: - filename = '/{}_{}.csv'.format(dbname, tablename) - - headers = [col[0] for col in cursor.description] # get headers - data.insert(0, tuple(headers)) - fp = open(concat(path,filename), 'w', newline = '') - myFile = csv.writer(fp) - myFile.writerows(data) - fp.close() - - except Exception as e: - print(e) - print("[FUNCTION_ERROR]-write_to_csv:" , filename) - - -#_______________MAIN_________________ - -#Change directory & open all file -change_directory_to(CURRENT_DIRECTORY_PATH) -MYSQL_CONFIG = read_json_file(MYSQL_CONFIG_FILE) -SQL_DICT = read_json_file(SQL_FILE) - - -#Connect DB -mydb = connect_mysql(MYSQL_CONFIG) -mycursor = mydb.cursor() - - -for tablename,sql in SQL_DICT.items(): - - try: - mycursor.execute(sql) - myresult = mycursor.fetchall() - dbname = mydb.database - write_to_csv(CSV_STORAGE_PATH,dbname,tablename,mycursor,myresult) - except Exception as e: - print(e) - print("[ERROR]-SQL ERROR",tablename,sql) - - -mydb.close() - - - diff --git a/python_el/extract/exchange_load/data-management-359208-6c514887e773.json b/python_el/extract/exchange_load/data-management-359208-6c514887e773.json deleted file mode 100644 index ee0d45c..0000000 --- a/python_el/extract/exchange_load/data-management-359208-6c514887e773.json +++ /dev/null @@ -1,12 +0,0 @@ -{ - "type": "service_account", - "project_id": "data-management-359208", - "private_key_id": "6c514887e773d5febd945655b02bb56396623933", - "private_key": "-----BEGIN PRIVATE KEY-----\nMIIEvwIBADANBgkqhkiG9w0BAQEFAASCBKkwggSlAgEAAoIBAQDKAgKEBSL06qKE\n9BJsnUndQrpUqzEEzeMVbFlOOP8G/Af6uBhnNp0fCU7C+bRwPV4Eyuyg36c3udn9\nMCUQbIT9fwF+/05NlmduPJg0DqrVVtzShHUv8lQTHHhKmIM9cmDKnVZ6vzVg2+pT\nE2P6289PR3FKZKZrZnsDNEbwIVfrSRqCQtfH4Kayqb1Y6IsGhog6NjbqFVyNN0pH\nlgxuHex/pHK2/U79rEq1tVmAJ1E4WNvJdZjt4MiHdkinvwD9pwAboZgxYnBDPynm\n2ElrYSzpvT4nLf8c0TpT0/dnd5BF/Eh15wyEcowdzR6xnqnfqKHOJzB+JTr7AKt8\nagr5UlvhAgMBAAECggEAEP7+SzFLcaPULK+EZVMOhek5WCpXI3pXItRM50HwYxwN\nZ9DZbMWxjozv7YOo5NCk+m5AXoCyxwOCDcVhOPKIdfObop3Ebs66wRGkFK0vPmfi\niGvQmEohPMJmdJBEaoUXE7UNM6Km0RFvs7Gr9c1MsfTm2UWCowKqUuixFz8W8Jq9\npV3VImRXr2P/Ftc2wrnabRIICgGIlzyuS3AAC+OI8QESavOyEh7Z9wGckRLL4YNq\nmNFfH9/PGq8qD9g7DEz+SC4VraTKga62AeadvSeRiYN9FuO5e32H6rag9D2d+0HM\n516SZjYPS7mkubI+7mbhSBhJRDk4bs0q5TxcFqyQ4QKBgQDxCNFvcvrdD4Kf8yO3\nD2wTwk1ukF/+tc2Kj1K84Z6IKuMTQtH7x1HgwoPssfuRjJLDwBkYELyAVgMSHh39\n3aw3KrxUw05wSSr4M4v282ykQcfThQjGkRgYBpZuRcVCqarAbC6eicM6lGots2QR\nofmUxBIX1Hc7eK8TQwBlXNlVjQKBgQDWjN+yDDd+/qEaeuz1rJsBSS6mjLSYlSYU\nZctd2oPPdc+bQEIeqkkjkqTc8R0gUXV9s5CwuAymyf9jhnVUQU//k+5epiP2wyUz\nb6vLGOeAhwpfqLNFZWuI/M2RioC1YaUAC2qJDnEJNoRBV8nXEsszV9JdEqzyOtba\n4AdAoV0YpQKBgQCZB/k4oi6l9XgAp3UQf6klrmJNBTr9U14JT8++/hwR5fC/xNfe\n3ACfC8CIocPP+AkiYS9NeSrE7FcMxLRT/s6dQ/PIeSuu3LV8WfXON2TNsLn3EGqu\n72X1sxEFOCTymxg/DTBYFa0u3xW+qDurekQkcIvwN0PwLUIyn4J72IRf7QKBgQCl\nwWBxZg7aBk7g7mdzxk5ax/dKpRpBZ7lruNlNQSzkcthZ0WND3btzyC+mooEmHsju\nvHPkk8zybszoT1EGLw9nHRrj9OeEFXAANR48Ypk4KxxQmz0lOB3ET8thzedyOmYH\nispb6NRbkcrL0M8XYmWq3Qag8XS8D8k+gCYaQJB0IQKBgQCIiwdoeS5c9M2IATaJ\nF8fKDSKKo44Zql6TwLOGFmxu/ENrPFDyz2OpbqQDCG1kWnY8dhtdCVNsUxDaNp99\nAKMJN2/lQCdZ7mdlm776jIVJi2lppYv+SuxvcKBNS1LcFnKQzhKZtXKJbzFW94Kw\noB/54kQof7RvtQdFGhhE2CIj4Q==\n-----END PRIVATE KEY-----\n", - "client_email": "exchange-python-el-bigquery@data-management-359208.iam.gserviceaccount.com", - "client_id": "102572291647345024418", - "auth_uri": "https://accounts.google.com/o/oauth2/auth", - "token_uri": "https://oauth2.googleapis.com/token", - "auth_provider_x509_cert_url": "https://www.googleapis.com/oauth2/v1/certs", - "client_x509_cert_url": "https://www.googleapis.com/robot/v1/metadata/x509/exchange-python-el-bigquery%40data-management-359208.iam.gserviceaccount.com" -} diff --git a/python_el/extract/exchange_load/exchange_load_main.py b/python_el/extract/exchange_load/exchange_load_main.py deleted file mode 100644 index faaac02..0000000 --- a/python_el/extract/exchange_load/exchange_load_main.py +++ /dev/null @@ -1,6 +0,0 @@ -import os -from google.cloud import bigquery -import json - -#Get all the Compulsory File -BIG_QUERY_CREDENTIAL = "data-management-359208-6c514887e773.json" diff --git a/run_extract_load.py b/run_extract_load.py new file mode 100644 index 0000000..98ac835 --- /dev/null +++ b/run_extract_load.py @@ -0,0 +1,54 @@ +import json +import os +from os import walk +from operator import concat +import general_function +import exchange.extract.exchange_extract_function as ex_extract_func +import exchange.load.exchange_load_function as ex_load_func + +def exchange_extract(EXCHANGE_CSV_TEMP_STORAGE_PATH): + #Config file name + EX_MYSQL_CONFIG_FILE = "EXCHANGE_CONFIG.json" + EX_SQL_FILE = "EXCHANGE_EXTRACT_SQL.json" + EXCHANGE_CSV_TEMP_STORAGE_PATH = EXCHANGE_CSV_TEMP_STORAGE_PATH + + #Read config file + EX_MYSQL_CONFIG = general_function.read_json_file("config/",EX_MYSQL_CONFIG_FILE) + EX_SQL_DICT = general_function.read_json_file("config/",EX_SQL_FILE) + + #Connect to mysql + EX_MYSQLDB = ex_extract_func.connect_mysql(EX_MYSQL_CONFIG) + EX_MYSQLCURSOR = EX_MYSQLDB.cursor() + + #Run EX_SQL_DICT to get all data from mysql in csv format + ex_extract_func.mysql_dict_sql_to_csv(EX_SQL_DICT, EXCHANGE_CSV_TEMP_STORAGE_PATH, EX_MYSQLDB, EX_MYSQLCURSOR) + +def exchange_load_to_cloudstorage(EXCHANGE_CSV_TEMP_STORAGE_PATH): + #Config file name + CLOUDSTORAGE_CREDENTIAL_FILE = "config/CLOUD-STORAGE-ADMIN-CREDENTIAL.json" + + #Connect to cloudstorage + ex_storage_client = ex_load_func.login_cloudstorage_credential(CLOUDSTORAGE_CREDENTIAL_FILE) + + temp_file_list = [] + for (dirpath, dirnames, filenames) in walk(EXCHANGE_CSV_TEMP_STORAGE_PATH): + temp_file_list.extend(filenames) + + if temp_file_list == []: + print("No Exchange CSV Found") + + for file_name in temp_file_list: + file_path = os.path.join(dirpath,file_name) + ex_load_func.upload_to_bucket(ex_storage_client, file_name, file_path) + + +def main(): + #Extract Load Exchange Data + EXCHANGE_CSV_TEMP_STORAGE_PATH = "exchange/extract/exchange_temp_extract_csv/" + exchange_extract(EXCHANGE_CSV_TEMP_STORAGE_PATH) + exchange_load_to_cloudstorage(EXCHANGE_CSV_TEMP_STORAGE_PATH) + + + +if __name__ == "__main__": + main() \ No newline at end of file