Prerequisite:
---------------------
AWS Lambda Layers Python | Snowflake-lambda-layer
• AWS Lambda Layers Python | Snowflake-...
Encrypt and Decrypt AWS Lambda Function Environment Variables using AWS KMS
• Encrypt and Decrypt AWS Lambda Functi...
Query Submitter Code:
---------------------------------------
import snowflake.connector as sf
import os
def run_query(conn, query):
print("Executing the query : ",query)
cursor = conn.cursor()
cursor.execute(query)
cursor.close()
def run_query1(conn, query):
print("Executing the query : ",query)
cursor = conn.cursor()
cursor.execute(query)
records=cursor.fetchone()
cursor.close()
return records
def execute_snowflake_query_async(conn,query):
cursor = conn.cursor()
cursor.execute_async(query)
query_id = cursor.sfqid
cursor.close()
print("Query ID : ",query_id)
return query_id
def lambda_handler(event, context):
print(event)
execute_the_query=event['queryStringParameters']['query']
print("Query to be executed : ",execute_the_query)
user=os.environ['username']
print("Username : ",user)
password=os.environ['password']
print("Password : ",password)
account="zk09286.us-east-2.aws";
database="DEMO_DB"
warehouse="COMPUTE_WH"
schema="PUBLIC"
role="ACCOUNTADMIN"
conn=sf.connect(user=user,password=password,account=account)
print("Connection successfully created")
statement_1='use warehouse '+warehouse;
statement3="use database "+database;
statement4="use role "+role;
run_query(conn,statement_1)
run_query(conn,statement3)
run_query(conn,statement4)
query_id=execute_snowflake_query_async(conn,execute_the_query)
return query_id
Query Poller Code:
----------------------------
import snowflake.connector as sf
import os
def run_query(conn, query):
cursor = conn.cursor()
cursor.execute(query)
cursor.close()
def run_query1(conn, query):
cursor = conn.cursor()
cursor.execute(query)
records=cursor.fetchone()
cursor.close()
return records
def get_result_from_query(conn,query_id):
cursor = conn.cursor()
cursor.get_results_from_sfqid(query_id)
records=cursor.fetchone()
cursor.close()
return records
def lambda_handler(event, context):
query_id=event['queryStringParameters']['query_id']
print("Query for which status has to be checked: ",query_id)
user=os.environ['username']
password=os.environ['password']
account="zk09286.us-east-2.aws";
database="DEMO_DB"
warehouse="COMPUTE_WH"
schema="PUBLIC"
role="ACCOUNTADMIN"
conn=sf.connect(user=user,password=password,account=account)
statement_1='use warehouse '+warehouse;
statement3="use database "+database;
statement4="use role "+role;
run_query(conn,statement_1)
run_query(conn,statement3)
run_query(conn,statement4)
query_to_be_executed="select EXECUTION_STATUS from table(information_schema.query_history()) where QUERY_ID='{}'".format(query_id)
status= run_query1(conn,query_to_be_executed)[0]
if(status=='SUCCESS'):
result=get_result_from_query(conn,query_id)[0]
return result
else:
return status
Reference:
-----------------
Perform Asynchronous Query Snowflake documentation --
https://docs.snowflake.com/en/user-gu...
Image of Architecture is taken from this blog --
https://streamhub.co.uk/an-approach-t...
Check this playlist for more AWS Projects in Big Data domain:
• Demystifying Data Engineering with Cl...
Auf dieser Seite können Sie das Online-Video Asynchronous layer over Snowflake | Execute Snowflake query with API mit der Dauer stunde minuten sekunde in guter Qualität ansehen, das der Benutzer Knowledge Amplifier 10 Januar 2022 hochgeladen hat, den Link mit Freunden und Bekannten teilen, dieses Video wurde auf Youtube bereits 1,302 Mal angesehen und es wurde von 25 den Zuschauern gefallen. Viel Spaß beim Betrachtenden Zuschauern gefallen!