Мне нужна помощь в скрипте PythonPython

Программы на Python
Anonymous
Мне нужна помощь в скрипте Python

Сообщение Anonymous »


Что сегодня (Фактически записи для обработки) 1. Подберите записи с END_DATE=31 декабря 9999. (Деактивировать текущую запись)2- Обновить дату окончания текущей записи на mis_date-1. (Новая запись) 3. Дата начала по умолчанию = Mis_date и end_date = 31 декабря 9999

Что предстоит разработать (Восстановление) 1- Выберите записи, в которых неправильная дата находится между начальной_датой и конечной_датой. 2- Обновить дату окончания текущей записи на mis_date-1. 3- Default_Start_Date = Mis_date и end_date = наименьшая будущая Start_Date -1
#!/usr/bin/python импортировать ОС импортировать систему импортировать панд как pd импортировать numpy как np из даты и времени импорта даты и времени, timedelta импортировать pyspark из столбца импорта pyspark.sql.functions из pyspark.sql.functions импортируйте max как max_ из pyspark импортировать SparkContext, SparkConf из импорта pyspark.sql.functions * из импорта pyspark.sql.functions горит из описания импорта pyspark.sql.functions из functools импортировать уменьшить из pyspark.sql импортировать DataFrame из окна импорта pyspark.sql.window из pyspark.sql импортировать SparkSession импортировать обратную трассировку журнал импорта из Connection.analyze_utility импортanalyze_table_partitioned,analyze_table_non_partitioned conf = SparkConf().setAppName("SPARK_INCREMENTALLOAD")\ .set('spark.kryoserializer.buffer.max','1024')\ .set('spark.kryoserializer.buffer','1024')\ .set('spark.executor.cores', 5)\ .set('spark.executor.memory', '25g')\ .set('spark.driver.memory', '16g')\ .set('spark.yarn.executor.memoryOverhead', '8192')\ .set('spark.yarn.queue', 'labstream')\ .set('spark.executor.heartbeatInterval','60s') искра = SparkSession.builder.enableHiveSupport().getOrCreate() spark.sparkContext.setLogLevel('ОШИБКА') spark.conf.set("hive.exec.dynamic.partition","true") spark.conf.set("hive.exec.dynamic.partition.mode","nonstrict") Ошибка класса (Исключение): """Базовый класс для других исключений""" проходить класс InvalidLoad (Ошибка): """Поднимается при попытке загрузки с обратной датой""" проходить print("Захват аргументов, переданных во время выполнения") пакетный_ид = sys.argv[1] mis_date = sys.argv[2] hive_db_name = sys.argv[3] hive_db_location = sys.argv[4] src_sys = sys.argv[5] src_stg_tab = sys.argv[6] пытаться: print("Чтение таблицы добавочной нагрузки установки"); setup_incremental_load = spark.sql("select * from " + hive_db_name + ".setup_incremental_load где source_table = '" + src_stg_tab.upper() + "' "); setup_incremental_load.cache() setup_incremental_load = setup_incremental_load.alias('setup_incremental_load') setup_incremental_load = setup_incremental_load.toPandas() src_stg_tab = src_stg_tab.lower() print("Имя исходной таблицы :: " + src_stg_tab) tgt_fct_tab = setup_incremental_load.target_table.iloc[0] tgt_fct_tab = tgt_fct_tab.lower() print("Имя целевой таблицы :: " + tgt_fct_tab) ##mis_date_in_date = datetime.strptime(mis_date, '%Y%m%d').strftime('%m/%d/%Y') ##mis_date_in_date = datetime.strptime(mis_date_in_date, '%m/%d/%Y').date() mis_date_in_date = datetime.strptime(str(mis_date),'%Y%m%d').strftime('%Y-%m-%d') Analysis_table_partitioned(hive_db_name,src_stg_tab,'d_mis_date='+ "'" + str(mis_date_in_date) + "'") анализ_таблицы_non_partitioned(hive_db_name,tgt_fct_tab) end_date = datetime.strptime('99991231', '%Y%m%d').strftime('%m/%d/%Y') end_date = datetime.strptime(end_date, '%m/%d/%Y').date() yes_flag="'" + 'Y' + "'" #hive_warehouse_path="/user/hive/warehouse" target_file_path=hive_db_location + "/" + tgt_fct_tab target_file_path_tmp=hive_db_location + "/" + "tmp_" + tgt_fct_tab print("Дата извлечения — изменение формата") mis_date = datetime.strptime(mis_date, '%Y%m%d').strftime('%Y-%m-%d') печать (неправильная_дата) prev_cnt=spark.sql("выберите d_mis_date из " + hive_db_name + ".com_load_history где v_target_tbl_name = '" + tgt_fct_tab.upper() + "' и \ v_load_status = '0' и d_mis_date > '" + mis_date_in_date+"'") prev_cnt_val=prev_cnt.count() если (prev_cnt_val > 0): поднять InvalidLoad ##Фиксация названий столбцов даты вступления в силу и окончания print("Столбец ключа даты вступления в силу"); eff_date = setup_incremental_load.target_column[(setup_incremental_load.source_table==src_stg_tab.upper()) & (setup_incremental_load.mapping_type=='EFFD')] eff_date = ''.join(eff_date.tolist()).lower() печать (eff_date) print("Столбец даты окончания"); end_date = setup_incremental_load.target_column[(setup_incremental_load.source_table==src_stg_tab.upper()) & (setup_incremental_load.mapping_type=='ED')] end_date = ''.join(end_date.tolist()).lower() печать(конечная_дата) print("Извлечение исходного системного ключа") src_sys_skey = spark.sql("выберите слияние(max(CAST(n_src_skey as int)),0) из " + hive_db_name + ".dim_source где v_src_code = '" + src_sys + "' и f_latest_record_indicator = " + yes_flag).limit( 1).collect()[0][0] печать (src_sys_skey) ## Создание кадра данных Spark как в исходной, так и в целевой таблицах src_tab = spark.sql("select * from " + hive_db_name + "." + src_stg_tab + " где d_mis_date = '" + mis_date + "' и v_src_code = '" + src_sys +"'") src_tab = src_tab.alias(src_stg_tab) src_tab.createOrReplaceTempView(src_stg_tab) src_tab_count=src_tab.count() print("\nПодсчет промежуточной таблицы ::" + str(src_tab_count)) если (src_tab_count == 0): print("\nНеверный запуск. В таблице этапов нет данных.\n") sys.exit(0) fct_full = spark.sql("выберите * из " + hive_db_name + "." + tgt_fct_tab) #fct_tab_closed_rec = spark.sql("select * from " + hive_db_name + "." + tgt_fct_tab + " где " + end_date + " 99991231 и n_src_skey = " + str(src_sys_skey)) #fct_tab_closed_rec = fct_tab_closed_rec.alias('fct_tab_closed_rec') #print("\nПодсчет ранее закрытых записей таблицы фактов::" + str(fct_tab_closed_rec.count())) #fct_tab = spark.sql("select * from " + hive_db_name + "." + tgt_fct_tab + " где " + end_date + " = 99991231 и n_src_skey = " + str(src_sys_skey) ) #fct_tab = fct_tab.alias('fct_tab') #fct_full = fct_tab_closed_rec.unionAll(fct_tab) #fct_tab.createOrReplaceTempView(tgt_fct_tab) #print("\nПодсчет количества активных записей таблицы фактов :: " + str(fct_tab.count())) #print("\nПодсчет количества закрытых записей таблицы фактов :: " + str(fct_tab_closed_rec.count())) fct_tbl_count = fct_full.count() если (fct_tbl_count == 0): пытаться: tmp_tbl_count = spark.sql("выберите счетчик(1) из " + hive_db_name + "."+"tmp_"+tgt_fct_tab).limit(1).collect()[0][0] #tmp_tbl_count = tmp_tbl_df_spk.count() кроме исключения как e: tmp_tbl_count = 0 если tmp_tbl_count >0: print("\nНеверный запуск. В таблице фактов нет данных, но есть данные в таблице tmp. Пожалуйста, свяжитесь со службой поддержки BAU\n"). sys.exit(50) fct_full.write.saveAsTable(hive_db_name+"."+"tmp_"+tgt_fct_tab,mode="overwrite",path=target_file_path_tmp) print("Данные загружаются в таблицу Temp") fct_tab_closed_rec = spark.sql("select * from " + hive_db_name + "." + "tmp_" + tgt_fct_tab + " где " + end_date + " 99991231 и n_src_skey = " + str(src_sys_skey)) fct_tab_closed_rec = fct_tab_closed_rec.alias('fct_tab_closed_rec') print("\nПодсчет ранее закрытых записей таблицы фактов::" + str(fct_tab_closed_rec.count())) fct_tab = spark.sql("select * from " + hive_db_name + "." + "tmp_" + tgt_fct_tab + " где " + end_date + " = 99991231 и n_src_skey = " + str(src_sys_skey) ) fct_tab = fct_tab.alias('fct_tab') #fct_full = fct_tab_closed_rec.unionAll(fct_tab) fct_tab.createOrReplaceTempView(tgt_fct_tab) print("\nПодсчет количества активных записей таблицы фактов :: " + str(fct_tab.count())) print("\nПодсчет количества закрытых записей таблицы фактов :: " + str(fct_tab_closed_rec.count())) ##Захват первичного ключа таблицы фактов без даты вступления в силу fct_pk_cols = setup_incremental_load.target_column[(setup_incremental_load.target_column_type=='PK') & (setup_incremental_load.mapping_type != 'EFFD')] fct_pk_cols = fct_pk_cols.tolist() print("Первичные ключи фактов :: "); fct_pk_cols = [element.lower() для элемента в fct_pk_cols] печать (fct_pk_cols) print("====================================Настройка общих переменных :: Начато === ========================================") ##Всего 7 переменных ##Идентификация столбцов прямого сопоставления таблицы фактов ##Переменная :: 1 Direct_map_cols = setup_incremental_load.target_column[(setup_incremental_load.mapping_type == 'DIRECT')] fct_direct_map_cols = [tgt_fct_tab + '.' + x вместо x в Direct_map_cols] fct_direct_map_cols = ','.join(fct_direct_map_cols) если fct_direct_map_cols: распечатать («не пусто») еще: fct_direct_map_cols = '-999' print("\nСтолбцы фактической прямой карты :: " + fct_direct_map_cols) ##Переменная :: 2 stg_direct_map_cols = [src_stg_tab + '.' + x вместо x в Direct_map_cols] stg_direct_map_cols = ','.join(stg_direct_map_cols) если stg_direct_map_cols: распечатать («не пусто») еще: stg_direct_map_cols = '-999' print("\nСтолбцы Stage Direct Map :: " + stg_direct_map_cols) date_format="'" + 'ггггММдд' + "'" печать (формат_даты) ##Формирование ключевых столбцов даты ##Переменная :: 3 date_key_cols_df = setup_incremental_load[(setup_incremental_load['mapping_type'] == 'DATE-KEY')] date_key_cols_df = date_key_cols_df.groupby('source_table').apply(lambda x:list( ' INT(from_unixtime(unix_timestamp(' + x.source_column + ',' + date_format + '),' + date_format + ' )) as ' + x.target_column )) для индекса значение в date_key_cols_df.iteritems(): date_key_cols = ','.join(val) если date_key_cols_df.empty: date_key_cols = '-999' еще: распечатать («не пусто») print("\nИсходные ключевые столбцы даты:: " + date_key_cols) ##Таблица фактов Столбцы с датами ##Переменная :: 4 fct_date_key_cols = setup_incremental_load.target_column[(setup_incremental_load.mapping_type == 'DATE-KEY')] fct_date_key_cols = [tgt_fct_tab + '.' + x вместо x в fct_date_key_cols] fct_date_key_cols = ','.join(fct_date_key_cols) если fct_date_key_cols: распечатать («не пусто») еще: fct_date_key_cols = '-999' print("\nФактические ключевые столбцы даты :: " + fct_date_key_cols) ##Определить столбцы таблицы Dim ##Переменная :: 5 dim_skey_cols = setup_incremental_load.supporting_dim_table_alias[(setup_incremental_load.mapping_type == 'SKEY')].astype(str)+ '.' + setup_incremental_load.supporting_dim_sk_col[(setup_incremental_load.mapping_type == 'SKEY')].astype(str)+' AS '+setup_incremental_load.target_column[(setup_incremental_load.mapping_type == 'SKEY')].astype(str) dim_skey_cols = ','.join(dim_skey_cols.tolist()) если dim_skey_cols: распечатать («не пусто») еще: dim_skey_cols = '-999' print("\nРазмер столбцов таблицы :: " + dim_skey_cols) ##Переменная :: 6 fct_skey_cols = setup_incremental_load.target_column[(setup_incremental_load.mapping_type == 'SKEY')] fct_skey_cols = ','.join(fct_skey_cols.tolist()) если fct_skey_cols: распечатать («не пусто») еще: fct_skey_cols = '-999' print("\nСтолбцы фактов SKEY :: " + fct_skey_cols) ##Формирование предложения соединения с таблицей исходного кода ##Переменная :: 7 join_clause = setup_incremental_load[(setup_incremental_load['supporting_dim_table'].notnull())] join_clause = join_clause.groupby('source_table').apply(lambda x:list(x.join_type + ' ' + hive_db_name + '.' + x.supporting_dim_table + ' '+ x.supporting_dim_table_alias + ' on ' + ' nvl( ' + x.source_table + '.' + x.source_column + ',-999) = nvl(' + x.supporting_dim_table_alias + '.' + x.supporting_dim_code_col + ',-999) и ' + x.supporting_dim_table_alias + '.f_latest_record_indicator = ' + да_флаг)) для индекса val в join_clause.iteritems(): join_stmt = ' '.join(val) print("\nОператор соединения с таблицей исходного кода :: " + join_stmt) print("====================================Настройка общих переменных :: Завершено === ========================================") print("===================================Закрытый набор данных :: Начато ===== =============================================== =====") df_closed_fct_rec=spark.sql("select " + fct_direct_map_cols + "," + fct_date_key_cols + "," + fct_skey_cols + " from " + tgt_fct_tab \ + " минус select " + stg_direct_map_cols + "," + date_key_cols + "," + dim_skey_cols + " from "\ + src_stg_tab + " " + join_stmt ); df_closed_fct_rec = df_closed_fct_rec.toDF(*[c.lower() для c в df_closed_fct_rec.columns]) df_closed_fct_rec = df_closed_fct_rec.alias('df_closed_fct_rec') fct_tab = fct_tab.alias('fct_tab') df_closed_fct_rec=fct_tab.join(df_closed_fct_rec,fct_pk_cols,"inner").select('fct_tab.*').select(fct_tab.columns) ##Обновление даты окончания закрытых записей end_date_val = datetime.strptime(str(mis_date_in_date), '%Y-%m-%d').date() - timedelta(days=1) end_date_val = end_date_val.strftime('%Y%m%d') ##df_closed_fct_rec=df_closed_fct_rec.withColumn('n_rating_end_date_skey',lit(int(end_date_val))) df_closed_fct_rec=df_closed_fct_rec.withColumn(end_date,when (df_closed_fct_rec[end_date] == 99991231,lit(int(end_date_val))).иначе(df_closed_fct_rec[end_date])) print("Отображение закрытого набора данных"); df_closed_fct_rec.show(5) Error_Log_Data=[(mis_date,batch_id,'FN_INCREMENTAL_LOAD',batch_id,'Поиск закрытых клиентов и ролей в fct','','FN_INCREMENTAL_LOAD: закрытые клиенты/роли в Fct','Шаг 1','','', mis_date_in_date,datetime.now().strftime('%Y-%m-%d %H:%M:%S'),'')] Error_Log_Data=spark.createDataFrame(Error_Log_Data) Error_Log_Data.createOrReplaceTempView("Error_Log_Data") spark.sql("вставить в таблицу " + hive_db_name + ".error_log_messages select * from Error_Log_Data"); print("===================================Закрытый набор данных :: Завершено ===== =============================================== =====") print("===================================Общий набор данных :: Начато ===== =============================================== ======") fct_table=spark.sql("select " + fct_direct_map_cols + "," + fct_date_key_cols + "," + fct_skey_cols + " from " + tgt_fct_tab) stg_table=spark.sql("select " + stg_direct_map_cols + "," + date_key_cols + "," + dim_skey_cols + " from " + src_stg_tab + " " + join_stmt ); df_common_stg_fct = fct_table.intersect(stg_table) df_common_stg_fct=fct_tab.join(df_common_stg_fct,fct_pk_cols,"inner").select('fct_tab.*').select(fct_tab.columns) print("Показываем общий набор данных"); df_common_stg_fct.cache() df_common_stg_fct.show() Error_Log_Data=[(mis_date,batch_id,'FN_INCREMENTAL_LOAD',batch_id,'Нахождение общих клиентов, fct и stg','','FN_INCREMENTAL_LOAD: общие клиенты в фактах и ​​стадиях ','Шаг 2','','',mis_date_in_date ,datetime.now().strftime('%Y-%m-%d %H:%M:%S'),'')] Error_Log_Data=spark.createDataFrame(Error_Log_Data) Error_Log_Data.createOrReplaceTempView("Error_Log_Data") spark.sql("вставить в таблицу " + hive_db_name + ".error_log_messages select * from Error_Log_Data"); print("===================================Общий набор данных :: Завершено ===== =============================================== =====") print("===================================Новый набор данных записи :: Начато ==== =============================================== ====") df_new_stg_rec=spark.sql("select " + stg_direct_map_cols + "," + date_key_cols + "," + dim_skey_cols + " from " + src_stg_tab + " "\ + " " + join_stmt + " минус select " + \ fct_direct_map_cols + "," + fct_date_key_cols + "," + fct_skey_cols + " from " + tgt_fct_tab); df_new_stg_rec=df_new_stg_rec.withColumn(eff_date,lit(int(sys.argv[2])).cast('bigint')) df_new_stg_rec=df_new_stg_rec.withColumn(end_date,lit(int(99991231)).cast('bigint')) df_new_stg_rec=df_new_stg_rec.withColumn('n_mis_date_skey',lit(int(sys.argv[2])).cast('bigint')).select(fct_tab.columns) print("Показываем новый набор данных"); df_new_stg_rec.show(5) Error_Log_Data=[(mis_date,batch_id,'FN_INCREMENTAL_LOAD',batch_id,'Поиск новых клиентов в stg','','FN_INCREMENTAL_LOAD: новые клиенты в таблице Stg ','Шаг 3','','',mis_date_in_date,datetime .now().strftime('%Y-%m-%d %H:%M:%S'),'')] Error_Log_Data=spark.createDataFrame(Error_Log_Data) Error_Log_Data.createOrReplaceTempView("Error_Log_Data") spark.sql("вставить в таблицу " + hive_db_name + ".error_log_messages select * from Error_Log_Data"); print("===================================Новый набор данных записи :: Завершено ==== =============================================== ====") print("====================================Объединенный набор данных :: Начато ===== =============================================== =========") ##Объединение фреймов данных print("Окончательное слияние :: Начато") Final_Row_Merged_spk = fct_tab_closed_rec.unionAll(df_closed_fct_rec).unionAll(df_common_stg_fct).unionAll(df_new_stg_rec) Final_Row_Merged_spk.cache() fct_cnt=Final_Row_Merged_spk.count() print("Общее количество :: " + str(fct_cnt)) Final_Row_Merged_spk.show(10) #Final_Row_Merged_spk.orderBy(Final_Row_Merged_spk.columns).write.mode("перезаписать").parquet(target_file_path) #Final_Row_Merged_spk.write.saveAsTable(hive_db_name+"."+"tmp_"+tgt_fct_tab,mode="overwrite",path=target_file_path_tmp) Final_Row_Merged_spk.write.saveAsTable(hive_db_name+"."+tgt_fct_tab,mode="overwrite",path=target_file_path) #print("Данные загружаются во временную таблицу") #spark.sql("вставить таблицу перезаписи "+hive_db_name+"."+tgt_fct_tab+" select * from "+hive_db_name+"."+"tmp_"+tgt_fct_tab) print("Данные загружаются в целевую таблицу") com_load_Data=[('INCREMENTAL_LOAD.PY',src_stg_tab.upper(),src_tab_count,fct_cnt,'0','',datetime.now().strftime('%Y-%m-%d %H:%M: %S'),src_sys,datetime.now().strftime('%Y-%m-%d'),tgt_fct_tab.upper(),mis_date)] com_load_Data=spark.createDataFrame(com_load_Data) com_load_Data.createOrReplaceTempView("com_load_Data") spark.sql("вставить в таблицу " + hive_db_name + ".com_load_history раздел (v_target_tbl_name,d_mis_date) выберите * из com_load_Data") Error_Log_Data=[(mis_date,batch_id,'FN_INCREMENTAL_LOAD',batch_id,'Объединение всех записей','','FN_INCREMENTAL_LOAD : объединение всех записей в таблице фактов ','Шаг 4','','',mis_date_in_date, datetime.now().strftime('%Y-%m-%d %H:%M:%S'),'')] Error_Log_Data=spark.createDataFrame(Error_Log_Data) Error_Log_Data.createOrReplaceTempView("Error_Log_Data") spark.sql("вставить в таблицу " + hive_db_name + ".error_log_messages select * from Error_Log_Data"); print("===================================Объединенный набор данных :: Завершено ===== =============================================== =======") кроме ValueError: print('Неверная дата!') sys.exit(50) кроме InvalidLoad: print("Попытка загрузки данных с предыдущей неверной датой!") Error_Log_Data=[(mis_date,batch_id,'FN_INCREMENTAL_LOAD',batch_id,'Попытка загрузки данных с предыдущей неправильной датой','','Попытка загрузки данных с предыдущей неправильной датой','Шаг 1 ','','', mis_date_in_date,datetime.now().strftime('%Y-%m-%d %H:%M:%S'),'')] Error_Log_Data=spark.createDataFrame(Error_Log_Data) Error_Log_Data.createOrReplaceTempView("Error_Log_Data") spark.sql("вставить в таблицу " + hive_db_name + ".error_log_messages выберите * из Error_Log_Data") sys.exit(50) кроме исключения как e: print("Ошибка - ",sys.exc_info()[0],"произошла.") Error_Log_Data=[(mis_date,batch_id,'FN_INCREMENTAL_LOAD',batch_id,sys.exc_info()[0],'',sys.exc_info()[0],'Step 5','','',mis_date_in_date,datetime. now().strftime('%Y-%m-%d %H:%M:%S'),'')] Error_Log_Data=spark.createDataFrame(Error_Log_Data) Error_Log_Data.createOrReplaceTempView("Error_Log_Data") spark.sql("вставить в таблицу " + hive_db_name + ".error_log_messages выберите * из Error_Log_Data") logging.error(traceback.format_exc()) sys.exit(50)

Вернуться в «Python»