Exchange Extract Load Function

This commit is contained in:
Yam ZhengLim
2022-08-22 18:54:54 +08:00
parent 7c34c88d0b
commit 39e707ce46
13 changed files with 166 additions and 116 deletions
+2 -1
View File
@@ -1 +1,2 @@
*.csv
*.csv
__pycache__
+8 -1
View File
@@ -7,4 +7,11 @@
"orders" : "SELECT * FROM orders;",
"payments" : "SELECT * FROM payments;",
"productlines" : "SELECT * FROM productlines;",
"products" : "SELECT * FROM products;"
"products" : "SELECT * FROM products;"
{
"transactions" : "SELECT * FROM transactions;",
"bookings" : "SELECT * FROM bookings;",
"service_types" : "SELECT * FROM service_types;"
}
@@ -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"
}
+10
View File
@@ -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;"
}
@@ -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)
+23
View File
@@ -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
+12
View File
@@ -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
@@ -1,5 +0,0 @@
{
"transactions" : "SELECT * FROM transactions;",
"bookings" : "SELECT * FROM bookings;",
"service_types" : "SELECT * FROM service_types;"
}
@@ -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()
@@ -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"
}
@@ -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"
+54
View File
@@ -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()