이전 이야기에서 공통코드테이블 매핑 작업을 진행하겠다는
일을 하기전에 앞서..
드디어 실데이터가 들어왔다.
(실데이터가 너무 늦게 들어오는바람에 공통코드테이블 일을 먼저 진행하고 있었던 것이기에, 기존에 하려했던 기본테이블 일감을 다시 진행)
우리는 계열사 전체에게 데이터를 받아야했기에,
계열사별 다루고 있는 데이터 방식이 다르다보니..
실데이터를 받자마자 우리가 예상했던 것과 다르게 데이터의 형태들은 엉망 진창으로 들어오게 되었다.
그래서 결국은 기존에 하던 개발은 멈추고 데이터 검증진행부터 하였다.
T0인 Temp_Table이 없는 관계로
airflow에 먼저 적재진행을하고 에러코드를 확인하여
리눅스 환경에서 데이터 검증을 진행하는 방식으로 하였다.
(왜 이렇게 하는지 이유는 모르겠지만.. 그러라고 하니까..)
12개의 계열사중 절반도 들어오지않은 데이터를 가지고도 이렇게 많은 오류가 나는것을 보아
지금 사용하고 있는 코드로는 무리가 있다는 생각을 하게되었다.
그래서 팀원들은 오류를 잡기위해 코드를 좀더 보강하고
명확한 오류를위해 row 단위로 데이터를 적재하는 코드로 수정하기로했다.
import pandas as pd
import chardet
import jaydebeapi
# 파일 인코딩 감지
file_path = 'path_to_your_file.dat'
with open(file_path, 'rb') as f:
result = chardet.detect(f.read())
encoding = result['encoding']
print(f"Detected encoding: {encoding}")
# .dat 파일 읽기 (감지된 인코딩 사용, 헤더 없음, 구분자는 ^)
df = pd.read_csv(file_path, delimiter='^', encoding=encoding, header=None)
# 데이터 파일의 실제 컬럼명 지정
df.columns = [
'직원번호'
, '기타_컬럼1'
, '기타_컬럼2'
, '기타_컬럼3'
]
# Tibero DB 연결 설정
import jaydebeapi
connection = jaydebeapi.connect('com.tmax.tibero.jdbc.TbDriver',
'jdbc:tibero:thin:@host:port:service_name',
['your_username', 'your_password'], 'path_to_tibero_jdbc_jar/tibero7-jdbc.jar')
cursor = connection.cursor()
# A테이블의 기존 데이터 삭제
delete_query = "DELETE FROM A테이블"
cursor.execute(delete_query)
connection.commit()
# 기본 테이블의 데이터 조회
query_basic = "SELECT 직원번호, 직원번호2 FROM 기본테이블"
cursor.execute(query_basic)
basic_table_data = cursor.fetchall()
# DataFrame으로 변환
basic_df = pd.DataFrame(basic_table_data, columns=['직원번호', '직원번호2'])
# 직원번호 열의 데이터 타입을 문자열로 변환
df['직원번호'] = df['직원번호'].astype(str)
basic_df['직원번호'] = basic_df['직원번호'].astype(str)
# 직원번호를 기준으로 매핑하여 직원번호2 추가
df = df.merge(basic_df, on='직원번호', how='left')
# 기본 키 값 확인
if df['직원번호'].isnull().any() or df['직원번호2'].isnull().any():
raise ValueError("PK(직원번호, 직원번호2)에 NULL 값이 포함되어 있습니다. 기본 키 값을 확인하십시오.")
# A 테이블에 데이터 삽입 (기타 컬럼 추가)
columns = df.columns.tolist() # DataFrame의 모든 컬럼명을 가져옴
columns_str = ', '.join(columns)
values_placeholders = ', '.join(['?'] * len(columns))
insert_query = f"""
INSERT INTO A테이블 ({columns_str})
VALUES ({values_placeholders})
"""
for index, row in df.iterrows():
cursor.execute(insert_query, tuple(row))
# 커밋 및 연결 종료
connection.commit()
cursor.close()
connection.close()
import pandas as pd
import chardet
import jaydebeapi
import logging
# 로깅 설정
logging.basicConfig(level=logging.ERROR)
logger = logging.getLogger(__name__)
# 파일 인코딩 감지
file_path = 'path_to_your_file.dat'
try:
with open(file_path, 'rb') as f:
result = chardet.detect(f.read())
encoding = result['encoding']
print(f"Detected encoding: {encoding}")
except Exception as e:
logger.error(f"{e} : Fail to detect file encoding")
raise
# .dat 파일 읽기 (감지된 인코딩 사용, 헤더 없음, 구분자는 ^)
try:
df = pd.read_csv(file_path, delimiter='^', encoding=encoding, header=None)
df.columns = ['직원번호', '기타_컬럼1', '기타_컬럼2', '기타_컬럼3']
except Exception as e:
logger.error(f"{e} : Fail to read the data file")
raise
# Tibero DB 연결 설정
try:
connection = jaydebeapi.connect(
'com.tmax.tibero.jdbc.TbDriver',
'jdbc:tibero:thin:@host:port:service_name',
['your_username', 'your_password'],
'path_to_tibero_jdbc_jar/tibero7-jdbc.jar'
)
cursor = connection.cursor()
except Exception as e:
logger.error(f"{e} : Fail to connect to Tibero database")
raise
# A테이블의 기존 데이터 삭제
try:
delete_query = "DELETE FROM A테이블"
cursor.execute(delete_query)
connection.commit()
except Exception as e:
logger.error(f"{e} : Fail to delete existing data from A테이블")
raise
# 기본 테이블의 데이터 조회
try:
query_basic = "SELECT 직원번호, 직원번호2 FROM 기본테이블"
cursor.execute(query_basic)
basic_table_data = cursor.fetchall()
basic_df = pd.DataFrame(basic_table_data, columns=['직원번호', '직원번호2'])
except Exception as e:
logger.error(f"{e} : Fail to fetch data from 기본테이블")
raise
# 직원번호 열의 데이터 타입을 문자열로 변환
try:
df['직원번호'] = df['직원번호'].astype(str)
basic_df['직원번호'] = basic_df['직원번호'].astype(str)
df = df.merge(basic_df, on='직원번호', how='left')
except Exception as e:
logger.error(f"{e} : Fail to process 직원번호 mapping")
raise
# 기본 키 값 확인
try:
if df['직원번호'].isnull().any() or df['직원번호2'].isnull().any():
raise ValueError("PK(직원번호, 직원번호2)에 NULL 값이 포함되어 있습니다. 기본 키 값을 확인하십시오.")
except Exception as e:
logger.error(f"{e} : PK validation failed")
raise
# A 테이블에 데이터 삽입 (row 단위로 처리)
try:
for row in df.itertuples(index=False):
insert_query =
f"INSERT INTO A테이블 (직원번호,
기타_컬럼1,
기타_컬럼2,
기타_컬럼3,
직원번호2)
VALUES (?,
?,
?,
?,
?)"
cursor.execute(insert_query, tuple(row))
connection.commit() # 각 행 삽입 후 커밋
except Exception as e:
logger.error(f"{e} : Fail to insert row data into A테이블")
raise
# 연결 종료
try:
cursor.close()
connection.close()
except Exception as e:
logger.error(f"{e} : Fail to close database connection")
raise
이번에 데이터 검증을 진행하면서 느꼈던 문제점은
결국 리눅스 환경에서 적재하는데는 한계가 있어 이부분에 대해 문제를 말씀드리고
Tibero DB에 각자 계정을 따서 테이블을 만들어 적재해 검증을 하게되었다.
다음주는.. 제발 데이터가 우리가 원하는 형식으로 들어오길..!!!