Merge pull request #217 feature-NEWDWH2021-1063 into develop

This commit is contained in:
下田雅人 2023-07-05 11:12:10 +09:00
commit caf5f5bde6
11 changed files with 1113 additions and 8 deletions

View File

@ -11,3 +11,9 @@ ULTMARC_BACKUP_FOLDER=ultmarc
JSKULT_CONFIG_BUCKET=**********************
JSKULT_CONFIG_CALENDAR_FOLDER=jskult/calendar
JSKULT_CONFIG_CALENDAR_HOLIDAY_LIST_FILE_NAME=jskult_holiday_list.txt
# 連携データ抽出期間
SALES_LAUNDERING_EXTRACT_DATE_PERIOD=0
# 洗替対象テーブル名
SALES_LAUNDERING_TARGET_TABLE_NAME=src05.sales_lau
# 卸実績洗替で作成するデータの期間(年単位)
SALES_LAUNDERING_TARGET_YEAR_OFFSET=5

View File

@ -1,5 +1,7 @@
from src.batch.common.batch_context import BatchContext
from src.batch.laundering import create_inst_merge_for_laundering, emp_chg_inst_laundering, ult_ident_presc_laundering
from src.batch.laundering import (
create_inst_merge_for_laundering, emp_chg_inst_laundering,
ult_ident_presc_laundering, sales_results_laundering)
from src.batch.dcf_inst_merge import integrate_dcf_inst_merge
from src.logging.get_logger import get_logger
@ -23,6 +25,8 @@ def exec():
emp_chg_inst_laundering.exec()
# 納入先処方元マスタ洗替
ult_ident_presc_laundering.exec()
# 卸販売洗替
sales_results_laundering.exec()
# # 並列処理のテスト用コード
# import time

View File

@ -0,0 +1,167 @@
from src.batch.batch_functions import logging_sql
from src.db.database import Database
from src.error.exceptions import BatchOperationException
from src.logging.get_logger import get_logger
from src.system_var import environment
logger = get_logger('卸販売洗替')
def exec():
db = Database.get_instance(autocommit=True)
try:
db.connect()
logger.debug('処理開始')
# 卸販売実績テーブル(洗替後)過去5年以前のデータ削除
_call_sales_lau_delete(db)
# 卸販売実績テーブル(洗替後)作成
_call_sales_lau_upsert(db)
# 1:卸組織洗替
_call_whs_org_laundering(db)
# HCO施設コードの洗替
_update_sales_lau_from_vop_hco_merge_v(db)
# 4:メルク施設コードの洗替
_update_mst_inst_laundering(db)
logger.debug('処理終了')
except Exception as e:
raise BatchOperationException(e)
finally:
db.disconnect()
def _call_sales_lau_delete(db: Database):
# 卸販売実績テーブル(洗替後)過去5年以前のデータ削除
logger.info('sales_lau_delete(プロシージャ―) 開始')
db.execute(f"""
CALL src05.sales_lau_delete(
'{environment.SALES_LAUNDERING_TARGET_TABLE_NAME}',
{environment.SALES_LAUNDERING_TARGET_YEAR_OFFSET}
)
""")
logger.info('sales_lau_delete(プロシージャ―) 終了')
return
def _call_sales_lau_upsert(db: Database):
# 卸販売実績テーブル(洗替後)作成
logger.info('sales_lau_upsert(プロシージャ―) 開始')
db.execute(f"""
CALL src05.sales_lau_upsert(
'{environment.SALES_LAUNDERING_TARGET_TABLE_NAME}',
(src05.get_syor_date() - interval {environment.SALES_LAUNDERING_EXTRACT_DATE_PERIOD} day),
src05.get_syor_date()
)
""")
logger.info('sales_lau_upsert(プロシージャ―) 終了')
return
def _call_whs_org_laundering(db: Database):
# 卸組織洗替
logger.info('whs_org_laundering(プロシージャ―) 開始')
db.execute(f"""
CALL src05.whs_org_laundering(
'{environment.SALES_LAUNDERING_TARGET_TABLE_NAME}'
)
""")
logger.info('whs_org_laundering(プロシージャ―) 終了')
return
def _update_sales_lau_from_vop_hco_merge_v(db: Database):
# HCO施設コードの洗替
if _count_v_inst_merge_t(db) == 0:
logger.info('V施設統合マスタ(洗替処理一時テーブル)にデータは存在しません')
return
_call_v_inst_merge_laundering(db)
return
def _count_v_inst_merge_t(db: Database) -> int:
# V施設統合マスタ(洗替処理一時テーブル)のデータ件数の取得
try:
sql = """
SELECT
COUNT(v_inst_cd) AS cnt
FROM
internal05.v_inst_merge_t
"""
result = db.execute_select(sql)
logging_sql(logger, sql)
logger.info('V施設統合マスタ(洗替処理一時テーブル)のデータ件数の取得 成功')
except Exception as e:
logger.debug('V施設統合マスタ(洗替処理一時テーブル)のデータ件数の取得 失敗')
raise e
return result[0]['cnt']
def _call_v_inst_merge_laundering(db: Database):
# HCO施設コードの洗替(プロシージャ―の呼び出し)
logger.info('v_inst_merge_laundering(プロシージャ―) 開始')
db.execute(f"""
CALL src05.v_inst_merge_laundering(
'{environment.SALES_LAUNDERING_TARGET_TABLE_NAME}'
)
""")
logger.info('v_inst_merge_laundering(プロシージャ―) 終了')
return
def _update_mst_inst_laundering(db: Database):
# メルク施設コードの洗替
_call_hco_to_mdb_laundering(db)
_update_sales_lau_from_dcf_inst_merge(db)
def _call_hco_to_mdb_laundering(db: Database):
# A:医療機関のデータはMDB変換表からHCO⇒DCFへ変換
logger.info('hco_to_mdb_laundering(プロシージャ―) 開始')
db.execute(f"""
CALL src05.hco_to_mdb_laundering(
'{environment.SALES_LAUNDERING_TARGET_TABLE_NAME}'
)
""")
logger.info('hco_to_mdb_laundering(プロシージャ―) 終了')
return
def _update_sales_lau_from_dcf_inst_merge(db: Database):
# B:DCF施設統合マスタがある場合は、コードを変換し、住所等をSETする
if _count_inst_merge_t(db) == 0:
logger.info('アルトマーク施設統合マスタ(洗替処理一時テーブル)にデータは存在しません')
return
_call_inst_merge_laundering(db)
return
def _count_inst_merge_t(db: Database) -> int:
# アルトマーク施設統合マスタ(洗替処理一時テーブル)のデータ件数の取得
try:
sql = """
SELECT
COUNT(dcf_dsf_inst_cd) AS cnt
FROM
internal05.inst_merge_t
"""
result = db.execute_select(sql)
logging_sql(logger, sql)
logger.info('アルトマーク施設統合マスタ(洗替処理一時テーブル)のデータ件数の取得 成功')
except Exception as e:
logger.debug('アルトマーク施設統合マスタ(洗替処理一時テーブル)のデータ件数の取得 失敗')
raise e
return result[0]['cnt']
def _call_inst_merge_laundering(db: Database):
# B:DCF施設統合マスタがある場合は、コードを変換し、住所等をSETする(プロシージャ―の呼び出し)
logger.info('inst_merge_laundering(プロシージャ―) 開始')
db.execute(f"""
CALL src05.inst_merge_laundering(
'{environment.SALES_LAUNDERING_TARGET_TABLE_NAME}'
)
""")
logger.info('inst_merge_laundering(プロシージャ―) 終了')
return

View File

@ -13,15 +13,17 @@ logger = get_logger(__name__)
class Database:
"""データベース操作クラス"""
__connection: Connection = None
__engine: Engine = None
__transactional_engine: Engine = None
__autocommit_engine: Engine = None
__host: str = None
__port: str = None
__username: str = None
__password: str = None
__schema: str = None
__autocommit: bool = None
__connection_string: str = None
def __init__(self, username: str, password: str, host: str, port: int, schema: str) -> None:
def __init__(self, username: str, password: str, host: str, port: int, schema: str, autocommit: bool = False) -> None:
"""このクラスの新たなインスタンスを初期化します
Args:
@ -30,12 +32,14 @@ class Database:
host (str): DBホスト名
port (int): DBポート
schema (str): DBスキーマ名
autocommit(bool): 自動コミットモードで接続するかどうか(Trueの場合トランザクションの有無に限らず即座にコミットされる). Defaults to False.
"""
self.__username = username
self.__password = password
self.__host = host
self.__port = int(port)
self.__schema = schema
self.__autocommit = autocommit
self.__connection_string = URL.create(
drivername='mysql+pymysql',
@ -47,16 +51,20 @@ class Database:
query={"charset": "utf8mb4"}
)
self.__engine = create_engine(
self.__transactional_engine = create_engine(
self.__connection_string,
pool_timeout=5,
poolclass=QueuePool
)
self.__autocommit_engine = self.__transactional_engine.execution_options(isolation_level='AUTOCOMMIT')
@classmethod
def get_instance(cls):
def get_instance(cls, autocommit=False):
"""インスタンスを取得します
Args:
autocommit (bool, optional): 自動コミットモードで接続するかどうか(Trueの場合トランザクションの有無に限らず即座にコミットされる). Defaults to False.
Returns:
Database: DB操作クラスインスタンス
"""
@ -65,7 +73,8 @@ class Database:
password=environment.DB_PASSWORD,
host=environment.DB_HOST,
port=environment.DB_PORT,
schema=environment.DB_SCHEMA
schema=environment.DB_SCHEMA,
autocommit=autocommit
)
@retry(
@ -77,12 +86,15 @@ class Database:
stop=stop_after_attempt(environment.DB_CONNECTION_MAX_RETRY_ATTEMPT))
def connect(self):
"""
DBに接続します接続に失敗した場合リトライします
DBに接続します接続に失敗した場合リトライします\n
インスタンスのautocommitがTrueの場合自動コミットモードで接続する明示的なトランザクションも無視される
Raises:
DBException: 接続失敗
"""
try:
self.__connection = self.__engine.connect()
self.__connection = (
self.__autocommit_engine.connect() if self.__autocommit is True
else self.__transactional_engine.connect())
except Exception as e:
raise DBException(e)

View File

@ -22,3 +22,10 @@ DB_CONNECTION_MAX_RETRY_ATTEMPT = int(os.environ.get('DB_CONNECTION_MAX_RETRY_AT
DB_CONNECTION_RETRY_INTERVAL_INIT = int(os.environ.get('DB_CONNECTION_RETRY_INTERVAL', 5))
DB_CONNECTION_RETRY_INTERVAL_MIN_SECONDS = int(os.environ.get('DB_CONNECTION_RETRY_MIN_SECONDS', 5))
DB_CONNECTION_RETRY_INTERVAL_MAX_SECONDS = int(os.environ.get('DB_CONNECTION_RETRY_MAX_SECONDS', 50))
# 連携データ抽出期間
SALES_LAUNDERING_EXTRACT_DATE_PERIOD = int(os.environ['SALES_LAUNDERING_EXTRACT_DATE_PERIOD'])
# 洗替対象テーブル名
SALES_LAUNDERING_TARGET_TABLE_NAME = os.environ['SALES_LAUNDERING_TARGET_TABLE_NAME']
# 卸実績洗替で作成するデータの期間(年単位)
SALES_LAUNDERING_TARGET_YEAR_OFFSET = os.environ['SALES_LAUNDERING_TARGET_YEAR_OFFSET']

View File

@ -0,0 +1,108 @@
-- A5M2で実行時に[SQL] - [スラッシュ(/)のみの行でSQLを区切る]に変えてから実行する
CREATE PROCEDURE src05.hco_to_mdb_laundering(target_table VARCHAR(255))
SQL SECURITY INVOKER
BEGIN
-- スキーマ名
DECLARE schema_name VARCHAR(50) DEFAULT (SELECT DATABASE());
-- プロシージャ名
DECLARE procedure_name VARCHAR(100) DEFAULT 'hco_to_mdb_laundering';
-- プロシージャの引数
DECLARE procedure_args JSON DEFAULT JSON_OBJECT('target_table', target_table);
-- 例外処理
DECLARE EXIT HANDLER FOR SQLEXCEPTION
BEGIN
GET DIAGNOSTICS CONDITION 1
@error_state = RETURNED_SQLSTATE, @error_msg = MESSAGE_TEXT;
CALL medaca_common.put_error_log(schema_name, procedure_name, procedure_args,
'hco_to_mdb_launderingでエラーが発生', @error_state, @error_msg);
SET @error_msg = (
CASE
WHEN LENGTH(@error_msg) > 128 THEN CONCAT(SUBSTRING(@error_msg, 1, 125), '...')
ELSE @error_msg
END
);
SIGNAL SQLSTATE '45000'
SET MYSQL_ERRNO = @error_state, MESSAGE_TEXT = @error_msg;
END;
SET @error_state = NULL, @error_msg = NULL;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】メルク施設コードの洗替_A① 開始');
TRUNCATE TABLE internal05.hco_cnv_mdb_t;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】メルク施設コードの洗替_A① 終了');
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】メルク施設コードの洗替_A② 開始');
INSERT INTO
internal05.hco_cnv_mdb_t (
hco_vid_v,
mdb_cd,
form_inst_name_kana,
form_inst_name_kanji,
inst_addr,
prefc_cd,
delete_flg,
abolish_ymd,
start_date
)
SELECT
mcmv.hco_vid_v,
mcmv.mdb_cd,
ci.form_inst_name_kana,
ci.form_inst_name_kanji,
ci.inst_addr,
ci.prefc_cd,
ci.delete_flg,
ci.abolish_ymd,
mcmv.start_date
FROM
src05.mdb_cnv_mst_v AS mcmv
INNER JOIN (
SELECT
hco_vid_v,MAX(sub_num) AS sno
FROM
src05.mdb_cnv_mst_v
WHERE
rec_sts_kbn != '9'
AND src05.get_syor_date() >= START_DATE
GROUP BY hco_vid_v
) AS mcmv2
ON mcmv.hco_vid_v = mcmv2.hco_vid_v
AND mcmv.sub_num = mcmv2.sno
LEFT OUTER JOIN src05.com_inst AS ci
ON mcmv.mdb_cd = ci.dcf_dsf_inst_cd
AND ci.delete_flg = '0'
;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】メルク施設コードの洗替_A② 終了');
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】メルク施設コードの洗替_A③ 開始');
SET @update_institution = "
UPDATE $$target_table$$ AS tt, internal05.hco_cnv_mdb_t AS hcmt
SET
tt.inst_cd = hcmt.mdb_cd,
tt.inst_name_kana = hcmt.form_inst_name_kana,
tt.inst_name = hcmt.form_inst_name_kanji,
tt.address = hcmt.inst_addr,
tt.pref_cd = hcmt.prefc_cd
WHERE
tt.v_inst_cd = hcmt.hco_vid_v
AND tt.inst_clas_cd = '1'
";
SET @update_institution = REPLACE(@update_institution, "$$target_table$$", target_table);
PREPARE update_institution_stmt from @update_institution;
EXECUTE update_institution_stmt;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】メルク施設コードの洗替_A③ 終了');
END

View File

@ -0,0 +1,65 @@
-- A5M2で実行時に[SQL] - [スラッシュ(/)のみの行でSQLを区切る]に変えてから実行する
CREATE PROCEDURE src05.inst_merge_laundering(target_table VARCHAR(255))
SQL SECURITY INVOKER
BEGIN
-- スキーマ名
DECLARE schema_name VARCHAR(50) DEFAULT (SELECT DATABASE());
-- プロシージャ名
DECLARE procedure_name VARCHAR(100) DEFAULT 'inst_merge_laundering';
-- プロシージャの引数
DECLARE procedure_args JSON DEFAULT JSON_OBJECT('target_table', target_table);
-- 例外処理
DECLARE EXIT HANDLER FOR SQLEXCEPTION
BEGIN
GET DIAGNOSTICS CONDITION 1
@error_state = RETURNED_SQLSTATE, @error_msg = MESSAGE_TEXT;
CALL medaca_common.put_error_log(schema_name, procedure_name, procedure_args,
'inst_merge_launderingでエラーが発生', @error_state, @error_msg);
SET @error_msg = (
CASE
WHEN LENGTH(@error_msg) > 128 THEN CONCAT(SUBSTRING(@error_msg, 1, 125), '...')
ELSE @error_msg
END
);
SIGNAL SQLSTATE '45000'
SET MYSQL_ERRNO = @error_state, MESSAGE_TEXT = @error_msg;
END;
SET @error_state = NULL, @error_msg = NULL;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】メルク施設コードの洗替_B① 開始');
SET @update_institution = "
UPDATE (
SELECT
dcf_dsf_inst_cd,
dup_opp_cd,
form_inst_name_kanji,
form_inst_name_kana,
inst_addr,
prefc_cd
FROM
internal05.inst_merge_t
) AS imt,
$$target_table$$ AS tt
SET
tt.inst_cd = imt.dup_opp_cd,
tt.inst_name = imt.form_inst_name_kanji,
tt.inst_name_kana = imt.form_inst_name_kana,
tt.address = imt.inst_addr,
tt.pref_cd = imt.prefc_cd,
tt.dwh_upd_dt = SYSDATE()
WHERE
tt.inst_cd = imt.dcf_dsf_inst_cd
AND tt.inst_clas_cd = '1'
";
SET @update_institution = REPLACE(@update_institution, "$$target_table$$", target_table);
PREPARE update_institution_stmt from @update_institution;
EXECUTE update_institution_stmt;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】メルク施設コードの洗替_B① 終了');
END

View File

@ -0,0 +1,49 @@
-- A5M2で実行時に[SQL] - [スラッシュ(/)のみの行でSQLを区切る]に変えてから実行する
CREATE PROCEDURE src05.sales_lau_delete(target_table VARCHAR(255), laundering_period_year INT)
SQL SECURITY INVOKER
BEGIN
-- スキーマ名
DECLARE schema_name VARCHAR(50) DEFAULT (SELECT DATABASE());
-- プロシージャ名
DECLARE procedure_name VARCHAR(100) DEFAULT 'sales_lau_delete';
-- プロシージャの引数
DECLARE procedure_args JSON DEFAULT JSON_OBJECT('target_table', target_table,
'laundering_period_year', laundering_period_year);
-- 例外処理
DECLARE EXIT HANDLER FOR SQLEXCEPTION
BEGIN
GET DIAGNOSTICS CONDITION 1
@error_state = RETURNED_SQLSTATE, @error_msg = MESSAGE_TEXT;
CALL medaca_common.put_error_log(schema_name, procedure_name, procedure_args,
'sales_lau_deleteでエラーが発生', @error_state, @error_msg);
SET @error_msg = (
CASE
WHEN LENGTH(@error_msg) > 128 THEN CONCAT(SUBSTRING(@error_msg, 1, 125), '...')
ELSE @error_msg
END
);
SIGNAL SQLSTATE '45000'
SET MYSQL_ERRNO = @error_state, MESSAGE_TEXT = @error_msg;
END;
SET @error_state = NULL, @error_msg = NULL;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'卸販売実績テーブル(洗替後)過去5年以前のデータ削除① 開始');
SET @delete_data = "
DELETE FROM
$$target_table$$
WHERE
kjyo_ym < DATE_FORMAT((src05.get_syor_date() - INTERVAL ? YEAR), '%Y%m')
";
SET @delete_data = REPLACE(@delete_data, "$$target_table$$", target_table);
PREPARE delete_data_stmt from @delete_data;
SET @interval_year = laundering_period_year;
EXECUTE delete_data_stmt USING @interval_year;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'卸販売実績テーブル(洗替後)過去5年以前のデータ削除① 終了');
END

View File

@ -0,0 +1,480 @@
-- A5M2で実行時に[SQL] - [スラッシュ(/)のみの行でSQLを区切る]に変えてから実行する
CREATE PROCEDURE src05.sales_lau_upsert(target_table VARCHAR(255), extract_from_date DATE,
extract_to_date DATE)
SQL SECURITY INVOKER
BEGIN
-- スキーマ名
DECLARE schema_name VARCHAR(50) DEFAULT (SELECT DATABASE());
-- プロシージャ名
DECLARE procedure_name VARCHAR(100) DEFAULT 'sales_lau_upsert';
-- プロシージャの引数
DECLARE procedure_args JSON DEFAULT JSON_OBJECT('target_table', target_table, 'extract_from_date',
extract_from_date, 'extract_to_date', extract_to_date);
-- 例外処理
DECLARE EXIT HANDLER FOR SQLEXCEPTION
BEGIN
GET DIAGNOSTICS CONDITION 1
@error_state = RETURNED_SQLSTATE, @error_msg = MESSAGE_TEXT;
CALL medaca_common.put_error_log(schema_name, procedure_name, procedure_args,
'sales_lau_upsertでエラーが発生', @error_state, @error_msg);
SET @error_msg = (
CASE
WHEN LENGTH(@error_msg) > 128 THEN CONCAT(SUBSTRING(@error_msg, 1, 125), '...')
ELSE @error_msg
END
);
SIGNAL SQLSTATE '45000'
SET MYSQL_ERRNO = @error_state, MESSAGE_TEXT = @error_msg;
END;
SET @error_state = NULL, @error_msg = NULL;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'卸販売実績テーブル(洗替後)作成① 開始'
);
TRUNCATE TABLE internal05.bu_prd_name_contrast_t;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'卸販売実績テーブル(洗替後)作成① 終了'
);
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'卸販売実績テーブル(洗替後)作成② 開始'
);
INSERT INTO
internal05.bu_prd_name_contrast_t (
prd_cd,
bu_cd,
phm_itm_cd,
pp_start_date,
pp_end_date,
update_date,
bp_start_date,
bp_end_date
)
SELECT
ppmv.prd_cd,
bpnc.bu_cd,
ppmv.phm_itm_cd,
ppmv.start_date AS pp_start_date,
ppmv.end_date AS pp_end_date,
bpnc.update_date AS update_date,
bpnc.start_date AS bp_start_date,
bpnc.end_date AS bp_end_date
FROM
src05.phm_prd_mst_v AS ppmv
LEFT OUTER JOIN src05.bu_prd_name_contrast AS bpnc
ON ppmv.phm_itm_cd = bpnc.phm_itm_cd
WHERE
ppmv.rec_sts_kbn != '9'
;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'卸販売実績テーブル(洗替後)作成② 終了'
);
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'卸販売実績テーブル(洗替後)作成③ 開始'
);
TRUNCATE TABLE internal05.fcl_mst_v_t;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'卸販売実績テーブル(洗替後)作成③ 終了'
);
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'卸販売実績テーブル(洗替後)作成④ 開始'
);
INSERT INTO
internal05.fcl_mst_v_t
SELECT
fmv1.v_inst_cd,
fmv1.sub_num,
fmv1.start_date,
fmv1.end_date,
fmv1.closed_dt,
fmv1.fcl_name,
fmv1.fcl_kn_name,
fmv1.fcl_abb_name,
fmv1.fcl_abb_kn_name,
fmv1.mkr_cd,
fmv1.jsk_proc_kbn,
fmv1.fmt_addr,
fmv1.fmt_kn_addr,
fmv1.postal_cd,
fmv1.prft_cd,
fmv1.prft_name,
fmv1.city_name,
fmv1.addr_line_1,
fmv1.tel_num,
fmv1.admin_kbn,
fmv1.fcl_type,
fmv1.rec_sts_kbn,
fmv1.ins_dt,
fmv1.upd_dt,
fmv1.dwh_upd_dt
FROM
src05.fcl_mst_v AS fmv1
INNER JOIN (
SELECT
fmv.v_inst_cd,
MAX(fmv.sub_num) AS sno
FROM
src05.fcl_mst_v AS fmv
GROUP BY
fmv.v_inst_cd
) AS fmv2
ON fmv1.v_inst_cd = fmv2.v_inst_cd
AND fmv1.sub_num = fmv2.sno
WHERE
fmv1.rec_sts_kbn != '9'
;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'卸販売実績テーブル(洗替後)作成④ 終了'
);
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'卸販売実績テーブル(洗替後)作成⑤ 開始'
);
SET @extract_from_datetime = CAST(extract_from_date AS DATETIME);
SET @extract_to_datetime = ADDTIME(CAST(extract_to_date AS DATETIME), '23:59:59');
SET @upsert_sales_launderning = "
INSERT INTO
$$target_table$$ (
rec_whs_cd,
rec_whs_sub_cd,
rec_whs_org_cd,
rec_cust_cd,
rec_comm_cd,
rec_tran_kbn,
rev_hsdnymd_wrk,
rev_hsdnymd_srk,
rec_urag_num,
rec_qty,
rec_nonyu_price,
rec_nonyu_amt,
rec_comm_name,
rec_nonyu_fcl_name,
free_item,
rec_nonyu_fcl_addr,
rec_nonyu_fcl_post,
rec_nonyu_fcl_tel,
rec_bef_hsdn_ymd,
rec_bef_slip_num,
rec_ymd,
sale_data_cat,
slip_file_name,
slip_mgt_num,
row_num,
hsdn_ymd,
exec_dt,
v_tran_cd,
tran_kbn_name,
whs_org_cd,
v_whsorg_cd,
whs_org_name,
whs_org_kn,
v_whs_cd,
whs_name,
nonyu_fcl_cd,
inst_name,
inst_name_kana,
address,
comm_cd,
comm_name,
nonyu_qty,
nonyu_price,
nonyu_amt,
shikiri_price,
shikiri_amt,
nhi_price,
nhi_amt,
v_inst_cd,
inst_clas_cd,
bu_cd,
item_cd,
item_name,
item_english_name,
pref_cd,
whspos_err_kbn,
htdnymd_err_kbn,
prd_exis_kbn,
fcl_exis_kbn,
bef_hsdn_ymd,
bef_slip_num,
slip_org_kbn,
kjyo_ym,
tksnbk_kbn,
fcl_exec_kbn,
rec_sts_kbn,
ins_dt,
ins_usr,
dwh_upd_dt
)
SELECT
s.rec_whs_cd,
s.rec_whs_sub_cd,
s.rec_whs_org_cd,
s.rec_cust_cd,
s.rec_comm_cd,
s.rec_tran_kbn,
s.rev_hsdnymd_wrk,
s.rev_hsdnymd_srk,
s.rec_urag_num,
s.rec_qty,
s.rec_nonyu_price,
s.rec_nonyu_amt,
s.rec_comm_name,
s.rec_nonyu_fcl_name,
s.free_item,
s.rec_nonyu_fcl_addr,
s.rec_nonyu_fcl_post,
s.rec_nonyu_fcl_tel,
s.rec_bef_hsdn_ymd,
s.rec_bef_slip_num,
s.rec_ymd,
s.sale_data_cat,
s.slip_file_name,
s.slip_mgt_num,
s.row_num,
s.hsdn_ymd,
s.exec_dt,
s.v_tran_cd,
s.tran_kbn_name,
s.whs_org_cd,
s.v_whsorg_cd,
s.whs_org_name,
s.whs_org_kn,
s.v_whs_cd,
s.whs_name,
s.nonyu_fcl_cd,
s.v_inst_name,
s.v_inst_kn,
s.v_inst_addr,
s.comm_cd,
s.comm_name,
CASE
WHEN
(LEFT(s.v_tran_cd, 1) = 2 AND (s.err_flg20 IS NULL OR s.err_flg20 != 'M'))
THEN
-s.nonyu_qty
ELSE
s.nonyu_qty
END AS nonyu_qty,
s.nonyu_price,
CASE
WHEN
(LEFT(s.v_tran_cd, 1) = 2 AND (s.err_flg20 IS NULL OR s.err_flg20 != 'M'))
THEN
-s.nonyu_amt
ELSE
s.nonyu_amt
END AS nonyu_amt,
s.shikiri_price,
CASE
WHEN
(LEFT(s.v_tran_cd, 1) = 2 AND (s.err_flg20 IS NULL OR s.err_flg20 != 'M'))
THEN
-s.shikiri_amt
ELSE
s.shikiri_amt
END AS shikiri_amt,
s.nhi_price,
CASE
WHEN
(LEFT(s.v_tran_cd,1) = 2 AND (s.err_flg20 IS NULL OR s.err_flg20 != 'M'))
THEN
-s.nhi_amt
ELSE
s.nhi_amt
END AS nhi_amt,
s.v_inst_cd,
CASE
WHEN
(fmvt.fcl_type = 'A1' or fmvt.fcl_type = 'A0') THEN '3'
WHEN
fmvt.fcl_type BETWEEN '20' AND '29' THEN '2'
ELSE
'1'
END AS inst_clas_cd,
bpnct.bu_cd,
ppmv.mkr_cd,
ppmv.mkr_inf_1,
ppmv.mkr_inf_2,
CASE
WHEN
s.v_inst_cd LIKE '00%'
THEN
ci.prefc_cd
ELSE
fmvt.prft_cd
END AS pref_cd,
s.whspos_err_kbn,
s.htdnymd_err_kbn,
s.prd_exis_kbn,
s.fcl_exis_kbn,
s.bef_hsdn_ymd,
s.bef_slip_num,
s.slip_org_kbn,
s.kjyo_ym,
s.tksnbk_kbn,
s.fcl_exec_kbn,
s.rec_sts_kbn,
s.ins_dt,
s.ins_usr,
SYSDATE()
FROM (
SELECT
? AS extract_from_datetime,
? AS extract_to_datetime
) AS sub
INNER JOIN src05.sales AS s
ON s.dwh_upd_dt BETWEEN sub.extract_from_datetime AND sub.extract_to_datetime
LEFT OUTER JOIN src05.phm_prd_mst_v AS ppmv
ON s.comm_cd = ppmv.prd_cd
AND STR_TO_DATE(s.hsdn_ymd,'%Y%m%d') BETWEEN ppmv.start_date AND ppmv.end_date
AND ppmv.rec_sts_kbn != '9'
LEFT OUTER JOIN internal05.fcl_mst_v_t AS fmvt
ON s.v_inst_cd = fmvt.v_inst_cd
LEFT OUTER JOIN internal05.bu_prd_name_contrast_t AS bpnct
ON s.comm_cd = bpnct.prd_cd
AND STR_TO_DATE(s.hsdn_ymd, '%Y%m%d') BETWEEN bpnct.pp_start_date AND bpnct.pp_end_date
AND STR_TO_DATE(s.hsdn_ymd, '%Y%m%d') BETWEEN bpnct.bp_start_date AND bpnct.bp_end_date
LEFT OUTER JOIN src05.com_inst AS ci
ON s.v_inst_cd = ci.dcf_dsf_inst_cd
WHERE
(s.rec_sts_kbn = '0' AND s.err_flg20 = 'M')
OR (
s.rec_sts_kbn = '0'
AND s.err_flg20 != 'M'
AND s.v_tran_cd IN (110, 120, 210, 220)
AND (
(s.fcl_exec_kbn NOT IN ('2', '5') AND (s.fcl_exec_kbn != '6' OR ppmv.prd_sale_kbn != 1))
OR s.fcl_exec_kbn IS NULL
)
)
ON DUPLICATE KEY UPDATE
rec_whs_cd = s.rec_whs_cd,
rec_whs_sub_cd = s.rec_whs_sub_cd,
rec_whs_org_cd = s.rec_whs_org_cd,
rec_cust_cd = s.rec_cust_cd,
rec_comm_cd = s.rec_comm_cd,
rec_tran_kbn = s.rec_tran_kbn,
rev_hsdnymd_wrk = s.rev_hsdnymd_wrk,
rev_hsdnymd_srk = s.rev_hsdnymd_srk,
rec_urag_num = s.rec_urag_num,
rec_qty = s.rec_qty,
rec_nonyu_price = s.rec_nonyu_price,
rec_nonyu_amt = s.rec_nonyu_amt,
rec_comm_name = s.rec_comm_name,
rec_nonyu_fcl_name = s.rec_nonyu_fcl_name,
free_item = s.free_item,
rec_nonyu_fcl_addr = s.rec_nonyu_fcl_addr,
rec_nonyu_fcl_post = s.rec_nonyu_fcl_post,
rec_nonyu_fcl_tel = s.rec_nonyu_fcl_tel,
rec_bef_hsdn_ymd = s.rec_bef_hsdn_ymd,
rec_bef_slip_num = s.rec_bef_slip_num,
rec_ymd = s.rec_ymd,
sale_data_cat = s.sale_data_cat,
slip_file_name = s.slip_file_name,
row_num = s.row_num,
hsdn_ymd = s.hsdn_ymd,
exec_dt = s.exec_dt,
v_tran_cd = s.v_tran_cd,
tran_kbn_name = s.tran_kbn_name,
whs_org_cd = s.whs_org_cd,
v_whsorg_cd = s.v_whsorg_cd,
whs_org_name = s.whs_org_name,
whs_org_kn = s.whs_org_kn,
v_whs_cd = s.v_whs_cd,
whs_name = s.whs_name,
nonyu_fcl_cd = s.nonyu_fcl_cd,
inst_name = s.v_inst_name,
inst_name_kana = s.v_inst_kn,
address = s.v_inst_addr,
comm_cd = s.comm_cd,
comm_name = s.comm_name,
nonyu_qty = VALUES(nonyu_qty),
nonyu_price = s.nonyu_price,
nonyu_amt = VALUES(nonyu_amt),
shikiri_price = s.shikiri_price,
shikiri_amt = VALUES(shikiri_amt),
nhi_price = s.nhi_price,
nhi_amt = VALUES(nhi_amt),
v_inst_cd = s.v_inst_cd,
inst_clas_cd = VALUES(inst_clas_cd),
bu_cd = bpnct.bu_cd,
item_cd = ppmv.mkr_cd,
item_name = ppmv.mkr_inf_1,
item_english_name = ppmv.mkr_inf_2,
pref_cd = VALUES(pref_cd),
whspos_err_kbn = s.whspos_err_kbn,
htdnymd_err_kbn = s.htdnymd_err_kbn,
prd_exis_kbn = s.prd_exis_kbn,
fcl_exis_kbn = s.fcl_exis_kbn,
bef_hsdn_ymd = s.bef_hsdn_ymd,
bef_slip_num = s.bef_slip_num,
slip_org_kbn = s.slip_org_kbn,
kjyo_ym = s.kjyo_ym,
tksnbk_kbn = s.tksnbk_kbn,
fcl_exec_kbn = s.fcl_exec_kbn,
rec_sts_kbn = s.rec_sts_kbn,
ins_dt = s.ins_dt,
ins_usr = s.ins_usr,
dwh_upd_dt = SYSDATE()
";
SET @upsert_sales_launderning = REPLACE(@upsert_sales_launderning, "$$target_table$$", target_table);
PREPARE upsert_sales_launderning_stmt from @upsert_sales_launderning;
EXECUTE upsert_sales_launderning_stmt USING @extract_from_datetime, @extract_to_datetime;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'卸販売実績テーブル(洗替後)作成⑤ 終了'
);
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'卸販売実績テーブル(洗替後)作成⑥ 開始'
);
SET @update_institution_code = "
UPDATE (
SELECT
? AS extract_from_datetime,
? AS extract_to_datetime
) AS sub,
$$target_table$$ AS tt,
src05.sales AS s
SET
tt.inst_cd = (
CASE
WHEN
(s.err_flg20 != 'M' AND tt.inst_clas_cd IN ('2', '3')) OR (s.err_flg20 = 'M')
THEN
s.v_inst_cd
ELSE
NULL
END
)
WHERE
s.dwh_upd_dt BETWEEN sub.extract_from_datetime AND sub.extract_to_datetime
AND tt.slip_mgt_num = s.slip_mgt_num
AND tt.row_num = s.row_num
";
SET @update_institution_code = REPLACE(@update_institution_code, "$$target_table$$", target_table);
PREPARE update_institution_code_stmt from @update_institution_code;
EXECUTE update_institution_code_stmt USING @extract_from_datetime, @extract_to_datetime;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'卸販売実績テーブル(洗替後)作成⑥ 終了'
);
END

View File

@ -0,0 +1,79 @@
-- A5M2で実行時に[SQL] - [スラッシュ(/)のみの行でSQLを区切る]に変えてから実行する
CREATE PROCEDURE src05.v_inst_merge_laundering(target_table VARCHAR(255))
SQL SECURITY INVOKER
BEGIN
-- スキーマ名
DECLARE schema_name VARCHAR(50) DEFAULT (SELECT DATABASE());
-- プロシージャ名
DECLARE procedure_name VARCHAR(100) DEFAULT 'v_inst_merge_laundering';
-- プロシージャの引数
DECLARE procedure_args JSON DEFAULT JSON_OBJECT('target_table', target_table);
-- 例外処理
DECLARE EXIT HANDLER FOR SQLEXCEPTION
BEGIN
GET DIAGNOSTICS CONDITION 1
@error_state = RETURNED_SQLSTATE, @error_msg = MESSAGE_TEXT;
CALL medaca_common.put_error_log(schema_name, procedure_name, procedure_args,
'v_inst_merge_launderingでエラーが発生', @error_state, @error_msg);
SET @error_msg = (
CASE
WHEN LENGTH(@error_msg) > 128 THEN CONCAT(SUBSTRING(@error_msg, 1, 125), '...')
ELSE @error_msg
END
);
SIGNAL SQLSTATE '45000'
SET MYSQL_ERRNO = @error_state, MESSAGE_TEXT = @error_msg;
END;
SET @error_state = NULL, @error_msg = NULL;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】HCO施設コードの洗替① 開始'
);
SET @update_institution = "
UPDATE (
SELECT
v_inst_cd,
v_inst_cd_merge,
fcl_name,
fcl_kn_name,
fmt_addr,
prft_cd
FROM
internal05.v_inst_merge_t
) AS vimt,
$$target_table$$ AS tt
SET
tt.inst_cd = (
CASE
WHEN
tt.inst_clas_cd = '1'
THEN
tt.inst_cd
WHEN
(tt.inst_clas_cd = '2' OR tt.inst_clas_cd = '3')
THEN
vimt.v_inst_cd_merge
END
),
tt.v_inst_cd = vimt.v_inst_cd_merge,
tt.inst_name = vimt.fcl_name,
tt.inst_name_kana = vimt.fcl_kn_name,
tt.address = vimt.fmt_addr,
tt.pref_cd = vimt.prft_cd,
tt.dwh_upd_dt = SYSDATE()
WHERE
tt.v_inst_cd = vimt.v_inst_cd
AND (tt.inst_clas_cd IN ('1', '2', '3'))
";
SET @update_institution = REPLACE(@update_institution, "$$target_table$$", target_table);
PREPARE update_institution_stmt from @update_institution;
EXECUTE update_institution_stmt;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】HCO施設コードの洗替① 終了'
);
END

View File

@ -0,0 +1,128 @@
-- A5M2で実行時に[SQL] - [スラッシュ(/)のみの行でSQLを区切る]に変えてから実行する
CREATE PROCEDURE src05.whs_org_laundering(target_table VARCHAR(255))
SQL SECURITY INVOKER
BEGIN
-- スキーマ名
DECLARE schema_name VARCHAR(50) DEFAULT (SELECT DATABASE());
-- プロシージャ名
DECLARE procedure_name VARCHAR(100) DEFAULT 'whs_org_laundering';
-- プロシージャの引数
DECLARE procedure_args JSON DEFAULT JSON_OBJECT('target_table', target_table);
-- 例外処理
DECLARE EXIT HANDLER FOR SQLEXCEPTION
BEGIN
GET DIAGNOSTICS CONDITION 1
@error_state = RETURNED_SQLSTATE, @error_msg = MESSAGE_TEXT;
CALL medaca_common.put_error_log(schema_name, procedure_name, procedure_args,
'whs_org_launderingでエラーが発生', @error_state, @error_msg);
SET @error_msg = (
CASE
WHEN LENGTH(@error_msg) > 128 THEN CONCAT(SUBSTRING(@error_msg, 1, 125), '...')
ELSE @error_msg
END
);
SIGNAL SQLSTATE '45000'
SET MYSQL_ERRNO = @error_state, MESSAGE_TEXT = @error_msg;
END;
SET @error_state = NULL, @error_msg = NULL;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】1.卸組織洗替① 開始'
);
TRUNCATE TABLE internal05.whs_customer_org_t;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】1.卸組織洗替① 終了'
);
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】1.卸組織洗替② 開始'
);
INSERT INTO
internal05.whs_customer_org_t (
whs_cd,
whs_sub_cd,
customer_cd,
whs_org_cd,
v_org_cd,
name_2
)
SELECT
wcmv.whs_cd,
wcmv.whs_sub_cd,
wcmv.customer_cd,
wcmv.whs_org_cd,
ocmv.v_org_cd,
mohv2.name_2
FROM
src05.whs_customer_mst_v AS wcmv
LEFT OUTER JOIN src05.org_cnv_mst_v AS ocmv
ON wcmv.whs_cd = ocmv.whs_cd
AND wcmv.whs_sub_cd = ocmv.whs_sub_cd
AND wcmv.whs_org_cd = ocmv.org_cd
AND src05.get_syor_date() BETWEEN ocmv.start_date AND ocmv.end_date
AND ocmv.rec_sts_kbn != '9'
LEFT OUTER JOIN (
SELECT
mohv.v_cd_2,
mohv.name_2
FROM src05.mkr_org_horizon_v AS mohv
INNER JOIN (
SELECT
v_cd_2,
MAX(dwh_upd_dt) AS dwh_upd_dt_latest
FROM
src05.mkr_org_horizon_v
WHERE
rec_sts_kbn != '9'
AND src05.get_syor_date() BETWEEN start_date AND end_date
GROUP BY
v_cd_2
ORDER BY
MAX(start_date) DESC
) AS m_latest
ON mohv.v_cd_2 = m_latest.v_cd_2
AND mohv.dwh_upd_dt = m_latest.dwh_upd_dt_latest
WHERE
mohv.rec_sts_kbn != '9'
AND src05.get_syor_date() BETWEEN mohv.start_date AND mohv.end_date
) AS mohv2
ON ocmv.v_org_cd = mohv2.v_cd_2
WHERE
wcmv.rec_sts_kbn != '9'
AND src05.get_syor_date() BETWEEN wcmv.start_date AND wcmv.end_date
;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】1.卸組織洗替② 終了'
);
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】1.卸組織洗替③ 開始'
);
SET @update_organization = "
UPDATE
$$target_table$$ AS tt, internal05.whs_customer_org_t AS wcot
SET
tt.whs_org_cd = wcot.whs_org_cd,
tt.v_whsorg_cd = wcot.v_org_cd,
tt.whs_org_name = wcot.name_2
WHERE
wcot.whs_cd = tt.rec_whs_cd
AND wcot.whs_sub_cd = tt.rec_whs_sub_cd
AND wcot.customer_cd = tt.rec_cust_cd
";
SET @update_organization = REPLACE(@update_organization, "$$target_table$$", target_table);
PREPARE update_organization_stmt from @update_organization;
EXECUTE update_organization_stmt;
CALL medaca_common.put_info_log(schema_name, procedure_name, procedure_args,
'【洗替】1.卸組織洗替③ 終了'
);
END