(
entity_df: Union[pd.DataFrame, str],
redshift_client,
config: RepoConfig,
s3_resource,
table_name: str,
)
| 1120 | |
| 1121 | |
| 1122 | def _upload_entity_df( |
| 1123 | entity_df: Union[pd.DataFrame, str], |
| 1124 | redshift_client, |
| 1125 | config: RepoConfig, |
| 1126 | s3_resource, |
| 1127 | table_name: str, |
| 1128 | ): |
| 1129 | if isinstance(entity_df, pd.DataFrame): |
| 1130 | # If the entity_df is a pandas dataframe, upload it to Redshift |
| 1131 | aws_utils.upload_df_to_redshift( |
| 1132 | redshift_client, |
| 1133 | config.offline_store.cluster_id, |
| 1134 | config.offline_store.workgroup, |
| 1135 | config.offline_store.database, |
| 1136 | config.offline_store.user, |
| 1137 | s3_resource, |
| 1138 | f"{config.offline_store.s3_staging_location}/entity_df/{table_name}.parquet", |
| 1139 | config.offline_store.iam_role, |
| 1140 | table_name, |
| 1141 | entity_df, |
| 1142 | ) |
| 1143 | elif isinstance(entity_df, str): |
| 1144 | # If the entity_df is a string (SQL query), create a Redshift table out of it |
| 1145 | aws_utils.execute_redshift_statement( |
| 1146 | redshift_client, |
| 1147 | config.offline_store.cluster_id, |
| 1148 | config.offline_store.workgroup, |
| 1149 | config.offline_store.database, |
| 1150 | config.offline_store.user, |
| 1151 | f"CREATE TABLE {table_name} AS ({entity_df})", |
| 1152 | ) |
| 1153 | else: |
| 1154 | raise InvalidEntityType(type(entity_df)) |
| 1155 | |
| 1156 | |
| 1157 | def _get_entity_schema( |
no test coverage detected