@@ -64,22 +64,30 @@ class QueryInformationOutput(DatabaseSettings, OutputSettings):
6464 DB_SOURCE_DESCRIPTION : str | None
6565
6666
67+ class QueryResultAsTable (DatabaseSettings , OutputSettings ):
68+ __identifier__ = "query_result_as_table"
69+
70+
6771class QueryDatabaseFromFileEntrypointSettings (EnvSettings ):
6872 DB_DSN : str
6973 DB_SCHEMA : str | None = None
7074
7175 query_file : QueryFileInput
76+
7277 csv_output : CSVOutput
7378 query_information : QueryInformationOutput
79+ query_result_as_table : QueryResultAsTable
7480
7581
7682class QueryDatabaseEntrypointSettings (EnvSettings ):
7783 DB_DSN : str
7884 DB_SCHEMA : str | None = None
7985
8086 query_str : QueryStrInput
87+
8188 csv_output : CSVOutput
8289 query_information : QueryInformationOutput
90+ query_result_as_table : QueryResultAsTable
8391
8492
8593def write_query_info (query : str , source : str , settings : DatabaseSettings ):
@@ -92,15 +100,22 @@ def write_query_info(query: str, source: str, settings: DatabaseSettings):
92100 db .write (table = settings .DB_TABLE , data = df , mode = "overwrite" )
93101
94102
95- @entrypoint (QueryDatabaseEntrypointSettings )
103+ def write_df_to_table (df : pd .DataFrame , settings : DatabaseSettings ) -> None :
104+ db = PandasDatabaseOperations (settings .DB_DSN , settings .DB_SCHEMA )
105+
106+ db .write (table = settings .DB_TABLE , data = df , mode = "overwrite" )
107+
108+
109+ # @entrypoint(QueryDatabaseEntrypointSettings)
96110def run_query_from_string (settings ):
97111 target_csv = "output.csv"
98- execute_query_to_csv (
112+ df = execute_query_to_csv (
99113 query = settings .query_str .QUERY ,
100114 dsn = settings .DB_DSN ,
101115 output_file = target_csv ,
102116 schema = settings .DB_SCHEMA ,
103117 )
118+ write_df_to_table (df , settings .query_result_as_table )
104119 upload_to_s3 (target_csv , settings .csv_output )
105120 write_query_info (
106121 query = settings .query_str .QUERY ,
@@ -122,12 +137,13 @@ def run_query_from_file(settings):
122137 query = read_query_file (local_file )
123138 target_csv = "output.csv"
124139
125- execute_query_to_csv (
140+ df = execute_query_to_csv (
126141 query = query ,
127142 dsn = settings .DB_DSN ,
128143 output_file = target_csv ,
129144 schema = settings .DB_SCHEMA ,
130145 )
146+ write_df_to_table (df , settings .query_result_as_table )
131147 upload_to_s3 (target_csv , settings .csv_output )
132148 write_query_info (
133149 query = query ,
0 commit comments