[SI] 프로젝트 20주차

소복치·2025년 2월 14일

오픈을 해야하는 시기이지만,
고객사 요청에 의해 오픈을 미루기로 했다.
그래서 그사이에 나는 코드를 좀 더 효율성 있게 고쳐보기로했다.

현재 내 코드로 50만건 정도 되는 데이터를 적재하려고 하면
기본 30분 정도 시간이 걸리는데
이것을 들으신 PM님 께서는 말이 안된다고 코드에 문제가 있는거 같다고 하셨다.

그래서 벌크로 들어가지만 row 단위로 들어갈때 처럼 에러를 잡을 수 있게 코드를 수정해 보았다.

import jaydebeapi as jp
import pandas as pd 
from datetime import datetime, timedelta
import errors as p
import glob
import os
import logging
import numpy as np
import shutil
import sys
from collections import defaultdict

logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')
logger = logging.getLogger(__name__)

def move_file(file, dir):
    base_name, extension = os.path.splitext(os.path.basename(file))
    destination = os.path.join(dir, os.path.basename(file))
    count = 1
    while os.path.exists(destination):
        new_file_name = f"{base_name}_{count}{extension}"
        destination = os.path.join(dir, new_file_name)
        count += 1
    shutil.move(file, destination)
    logger.info(f"파일이 {destination}으로 이동되었습니다.")

# JAR파일 경로
jar_path_local = 'jdbc 파일 경로'

# 데이터 파일 경로
dat_file_path = '파일 경로'
dat_file_name = '파일 이름'

# 성공 및 실패 파일 이동 경로
fail_dir = '실패 파일 경로'
success_dir = '성공 파일 경로'
os.makedirs(success_dir, exist_ok=True)
os.makedirs(fail_dir, exist_ok=True)

files = glob.glob(os.path.join(dat_file_path, dat_file_name))
if not files:
    logger.info("파일이 존재하지 않아 종료합니다.")
    sys.exit()

logger.info(f"Insert File Start :: {datetime.now().strftime('%Y%m%d%H%M%S')}")
conn = p.conn_jdbc()
conn.jconn.setAutoCommit(False)
cursor = conn.cursor()

due_time = datetime.now().strftime('%Y-%m-%d %H:%M:%S')
process_dt = (datetime.now() - timedelta(days=1)).strftime('%Y%m%d')
act_dt = datetime.now().strftime('%Y%m%d')

def merge_into_origin():
    global merge_chk
    try:
        cursor.execute("""
            MERGE INTO 테이블명 tgt
            USING 테이블명 src
            ON (tgt.직원번호 = src.직원번호)
            WHEN MATCHED THEN
                UPDATE SET
                    tgt.기타컬럼1 = src.기타컬럼1,
                    tgt.기타컬럼2 = src.기타컬럼2,
                    tgt.기타컬럼3 = src.기타컬럼3
            WHEN NOT MATCHED THEN
                INSERT (직원번호, 기타컬럼1, 기타컬럼2, 기타컬럼3)
                VALUES (src.직원번호, src.기타컬럼1, src.기타컬럼2, src.기타컬럼3);
        """)
        conn.commit()
        merge_chk = True
    except jp.DatabaseError as e:
        logger.error(f"Database error during merge: {e}")
        merge_chk = False
    except Exception as e:
        logger.error(f"Unexpected error during merge: {e}")
        merge_chk = False

try:
    for file in files:
        logger.info(f"Processing file: {os.path.basename(file)}")
        totalcnt, errorcnt, inserterrorcnt = 0, 0, 0
        GROUP_CO_CD = os.path.basename(file)[:3].upper()
        csv_data = p.read_csv(file)

        if csv_data is not None:
            totalcnt = len(csv_data)
            try:
                csv_data.columns = ['직원번호', '기타컬럼1', '기타컬럼2', '기타컬럼3']
            except Exception as e:
                logger.error(f"Column mapping error: {e}")
                errorcnt += 1

            csv_data = csv_data.apply(p.chk_chg_empid, axis=1)
            csv_data = csv_data.replace({'"': '', "'": ''}, regex=True)
            
            p.insert_log(cursor, 'start', '테이블명', GROUP_CO_CD, due_time, process_dt, act_dt)

            insert_sql = """
                INSERT INTO 테이블명 (직원번호, 기타컬럼1, 기타컬럼2, 기타컬럼3)
                VALUES (?, ?, ?, ?)
            """

            BATCH_SIZE = 1000
            total_inserted = 0
            for i in range(0, len(csv_data), BATCH_SIZE):
                batch = csv_data.iloc[i:i+BATCH_SIZE]
                try:
                    cursor.executemany(insert_sql, batch.values.tolist())
                    conn.commit()
                    total_inserted += len(batch)
                except Exception as e:
                    conn.rollback()
                    logger.error(f"Error inserting batch: {e}")

            merge_into_origin()
            fails = totalcnt - total_inserted
            p.insert_log(cursor, 'end', '테이블명', GROUP_CO_CD, due_time, process_dt, act_dt, totalcnt, total_inserted, fails)
            move_file(file, success_dir if merge_chk else fail_dir)
        else:
            errorcnt += 1
            move_file(file, fail_dir)

        logger.info(f"Processing complete: {os.path.basename(file)}")

finally:
    for file in glob.glob(os.path.join(dat_file_path, "*.chk")):
        try:
            os.remove(file)
            logger.info(f"Deleted: {file}")
        except Exception as e:
            logger.error(f"Error deleting {file}: {e}")

    try:
        cursor.close()
        conn.close()
    except Exception as e:
        logger.error(f"Error closing connection: {e}")

주요 코드 설명
1. 파일 이동 함수(move_file)

  • 파일 충돌 방지 : 같은 이름의 파이링 있으면 1,2를 붙여가며 저장
  • shutil.move() : 파일을 성공/실패 폴더로 이동
  1. 파일 존재 여부 확인
  • 파일이 없으면 프로그램 종료
  • glob.glob() 를 사용해 특정 패턴의 파일 목록 가져오기
  1. 데이터베이스 연결
  • JDBC 연결
  • 트랜잭션 관리 (setAutoCommit(False)) → commit() 또는 rollback()
  1. MERGE INTO
  • 기존 데이터 업데이트 & 새로운 데이터 삽입
  • WHEN MATCHED → UPDATE, WHEN NOT MATCHED → INSERT
  • 오류 발생 시 False 반환 -> rollback() 필요
  1. 데이터 읽기 및 처리
  • p.read_csv(file): CSV 파일 읽기
  • 컬럼명 설정 및 데이터 정제
  • p.chk_chg_empid() 적용 (사번 등 변경 사항 반영)
  1. 배치 단위 데이터 삽입
  • 배치 처리 (BATCH_SIZE = 1000) → 속도 최적화
  • executemany() → 여러 개의 데이터를 한 번에 INSERT
  • 삽입 실패 시 rollback()
  1. 파일 성공/실패 처리
  • merge_chk 값에 따라 성공/실패 폴더로 파일 이동
  • 최종 로그 기록
  1. 파일 정리 & 연결 종료
  • 임시 파일(.chk) 삭제
  • DB 커서 및 연결 종료 (cursor.close(), conn.close())

이로써 오픈까지 모든 준비는 끝났고,
이제 오픈하고 나서 생기는 에러들을 잘 잡으면 될거 같다.

첫 프로젝트인 만큼 실수도 많았고, 모르는것도 많았고 너무 힘들었지만
막상 프로젝트가 끝나가니 조금은 내가 성장했다는것을 내 스스로가 느꼈다.

다같이 함께한 팀원들과의 호흡도 너무 좋았고 프로젝트를 또 다시 한번 더
하고 싶다는 생각을 많이 하게됐다!

profile
오늘 터져도내일 다시극복

1개의 댓글

comment-user-thumbnail
2025년 2월 14일

우선, 프로젝트 연기 동안 코드를 최적화하려는 결심은 좋은 생각입니다. 데이터 적재 시간을 줄이고 에러를 더 효과적으로 처리하는 것은 매우 중요하죠. 여기 몇 가지 개선할 수 있는 팁을 드릴게요:
iTero

답글 달기