
안녕하세요 가야태자 @talkit 입니다.
@talkit.bank 서비스와 talkitsteem 프로젝트를 진행하고 있는데 이 프로그램들의 제어를 위해서 메시지 기반 처리 서비스를 생각하고 있습니다.
메시지 기반 처리 서비스 중에 백업을 먼저 고려하고 앞에서 말씀 드렸지만, Kafka 실행이 실패 했습니다. T.T
그래서 Python으로 필요한 데이터를 복사하는 작은 프로그램을 작성해 보았습니다.
정확히 말하면 이 코드는 Kafka도 CDC(Change Data Capture)도 아닙니다. 변경 로그를 지속해서 읽는 것이 아니라 실행 시점의 일부 데이터를 비교해 넣는 일회성 동기화 실험입니다. 삭제 전파, 체크포인트, 재시도, 순서 보장과 전체 데이터 페이지 처리가 없으므로 운영 복제 도구로 사용하면 안 됩니다.
아래의 오래된 실험 코드는 구조를 설명하기 위해 남겼습니다. 코드 안에 DB 주소·사용자·비밀번호를 직접 넣은 방식은 따라 하지 마십시오. 실제 구현에서는 환경 변수나 시크릿 저장소를 사용하고, 소스는 읽기 전용 계정, 대상은 필요한 테이블에만 쓰기 가능한 별도 계정을 사용해야 합니다. TLS 인증서 검증과 개인정보 마스킹도 필요합니다.
소스 MySQL 접속
타겟 MySQL 접속
소스에 있는 필요한 테이블 정의
타겟에 테이블이 존재한다고 생각
실제 필드에서 CDC를 진행한다면 처음 일정 부분의 데이터를 부어 놓고 일을 하겠지만, ^^ 저는 개발 서브를 만들기 위해서 시작 한거라 해당 작업은 넘어 갔습니다.
여러테이블에서 PK기준 데이터 검색
데이터가 없으면 Insert
데이터가 있으면 Update
데이터가 있으면 입력할 데이터와 전체 데이터 비교
이때 데이터 변경이 없으면 다음으로 넘어가고 있으면 Update를 수행 했습니다.
삭제 관련 고려 없음.
원래는 Delete도 고려 해야하는데 Delete는 고려하지 않았습니다.
import pymysql
import logging
# 로그 설정
logging.basicConfig(level=logging.INFO)
# 소스 MySQL에 연결하는 함수
def connect_mysql(host, user, password, database):
try:
connection = pymysql.connect(
host=host,
user=user,
password=password,
database=database,
charset='utf8'
)
return connection
except pymysql.MySQLError as e:
logging.error(f"MySQL connection error: {e}")
return None
# 소스 MySQL에서 데이터를 조회하는 함수
def fetch_data_from_source(connection, table_name):
cursor = connection.cursor(pymysql.cursors.DictCursor)
try:
# 원하는 테이블에서 1000건 조회
query = f"SELECT * FROM {table_name} LIMIT 1000"
cursor.execute(query)
data = cursor.fetchall()
return data
except pymysql.MySQLError as e:
logging.error(f"Error fetching data from source: {e}")
return None
finally:
cursor.close()
# 타겟 MySQL에서 PK가 존재하는지 조회하는 함수
def check_pk_in_target(connection, table_name, pk_value):
cursor = connection.cursor()
try:
if table_name == 'sample_asset_daily':
query = f"SELECT COUNT(*) FROM {table_name} WHERE collect_time = %s AND asset = %s AND product_id = %s"
elif table_name == 'sample_user_daily':
query = f"SELECT COUNT(*) FROM {table_name} WHERE user_id = %s AND year = %s AND month = %s AND day = %s"
elif table_name == 'sample_user_monthly':
query = f"SELECT COUNT(*) FROM {table_name} WHERE user_id = %s AND year = %s AND month = %s"
elif table_name == 'sample_users':
query = f"SELECT COUNT(*) FROM {table_name} WHERE user_id = %s"
else:
return False
cursor.execute(query, pk_value)
count = cursor.fetchone()[0]
return count > 0 # 존재하면 True, 없으면 False
except pymysql.MySQLError as e:
logging.error(f"Error checking PK in target: {e}")
return False
finally:
cursor.close()
# 타겟 MySQL에 데이터를 Insert 하는 함수
def insert_into_target(connection, table_name, data):
cursor = connection.cursor()
try:
columns = ', '.join(data.keys())
values = ', '.join(['%s'] * len(data))
query = f"INSERT INTO {table_name} ({columns}) VALUES ({values})"
cursor.execute(query, tuple(data.values()))
connection.commit()
logging.info(f"Inserted data into {table_name}")
except pymysql.MySQLError as e:
logging.error(f"Error inserting data into target: {e}")
connection.rollback()
finally:
cursor.close()
# 타겟 MySQL에 데이터를 Update 하는 함수
def update_target(connection, table_name, data, pk_value):
cursor = connection.cursor()
try:
set_clause = ', '.join([f"{key} = %s" for key in data.keys()])
if table_name == 'sample_asset_daily':
query = f"UPDATE {table_name} SET {set_clause} WHERE collect_time = %s AND asset = %s AND product_id = %s"
elif table_name == 'sample_user_daily':
query = f"UPDATE {table_name} SET {set_clause} WHERE user_id = %s AND year = %s AND month = %s AND day = %s"
elif table_name == 'sample_user_monthly':
query = f"UPDATE {table_name} SET {set_clause} WHERE user_id = %s AND year = %s AND month = %s"
elif table_name == 'sample_users':
query = f"UPDATE {table_name} SET {set_clause} WHERE user_id = %s"
cursor.execute(query, tuple(data.values()) + pk_value)
connection.commit()
logging.info(f"Updated data in {table_name}")
except pymysql.MySQLError as e:
logging.error(f"Error updating data in target: {e}")
connection.rollback()
finally:
cursor.close()
# 메인 로직
def process_data():
source_connection = connect_mysql('<SOURCE_DB_HOST>', '<SOURCE_DB_USER>', '<SOURCE_DB_PASSWORD>', '<SOURCE_DB_NAME>')
target_connection = connect_mysql('<TARGET_DB_HOST>', '<TARGET_DB_USER>', '<TARGET_DB_PASSWORD>', '<TARGET_DB_NAME>')
if not source_connection or not target_connection:
logging.error("Failed to connect to MySQL.")
return
# CDC 대상 테이블 리스트
target_tables = ['sample_asset_daily', 'sample_user_daily', 'sample_user_monthly', 'sample_users']
for table in target_tables:
# 소스 MySQL에서 데이터 조회
data = fetch_data_from_source(source_connection, table)
if not data:
logging.error(f"No data found in source table {table}")
continue
for row in data:
# 테이블별 PK 처리
if table == 'sample_asset_daily':
pk_value = (row['collect_time'], row['asset'], row['product_id'])
elif table == 'sample_user_daily':
pk_value = (row['user_id'], row['year'], row['month'], row['day'])
elif table == 'sample_user_monthly':
pk_value = (row['user_id'], row['year'], row['month'])
elif table == 'sample_users':
pk_value = (row['user_id'],)
# 타겟 MySQL에서 해당 PK가 있는지 확인
if check_pk_in_target(target_connection, table, pk_value):
# 타겟에서 해당 데이터가 있으면, 기존 데이터와 비교하여 다르면 Update
logging.info(f"PK found in target for {table}, checking data.")
if row != fetch_data_from_source(target_connection, table): # 데이터가 다르면 업데이트
update_target(target_connection, table, row, pk_value)
else:
logging.info(f"No changes for PK {pk_value} in table {table}. Skipping update.")
else:
# PK가 없다면 Insert
insert_into_target(target_connection, table, row)
source_connection.close()
target_connection.close()
if __name__ == '__main__':
process_data()소스는 위와 같습니다. 다만 현재 상태 그대로 운영에 사용해서는 안 됩니다.
이 예제에는 다음 한계가 있습니다.
LIMIT 1000이후 데이터를 읽지 않습니다- 삭제된 데이터가 대상에 반영되지 않습니다
- 중단 뒤 이어서 처리할 체크포인트와 재시도가 없습니다
- 기존 행 비교 코드가 한 행과 전체 조회 결과를 비교해 의도대로 동작하지 않습니다
- 테이블명을 문자열로 SQL에 넣으므로 외부 입력을 받지 않고 허용 목록으로 고정해야 합니다
- PK 값과 데이터 내용이 로그에 노출되지 않도록 해야 합니다
실제 접속 코드는 다음처럼 시크릿을 코드 밖에서 받아야 합니다.
import os
import pymysql
source_connection = pymysql.connect(
host=os.environ["SOURCE_DB_HOST"],
user=os.environ["SOURCE_DB_USER"],
password=os.environ["SOURCE_DB_PASSWORD"],
database=os.environ["SOURCE_DB_NAME"],
charset="utf8mb4",
connect_timeout=10,
ssl={"ca": os.environ["MYSQL_CA_FILE"]},
)
운영 DB를 개발 DB로 복사할 때는 계정·개인정보·토큰을 먼저 마스킹해야 합니다. 이 글의 코드는 실패한 첫 실험으로 읽어 주시고, 운영용 CDC는 MySQL binlog 기반 도구와 재처리 정책까지 포함해 별도로 설계해야 합니다.
해당 소스를 이용해서 저는 UI 개발에 박차를 가하겠습니다.
감사합니다.
Tags MySQL, Python, 데이터복사, CDC, 데이터동기화, 환경변수, TLS, 데이터마스킹, 개발실험, 가야태자, talkit
'DevOps·CI·CD' 카테고리의 다른 글
| [반자동 포스팅기] 3. 손으로 눌러도 안 됐습니다 (0) | 2026.09.28 |
|---|---|
| [반자동 포스팅기] 2. 제가 만든 도구가, 저에게 거짓말을 했습니다 (0) | 2026.09.21 |
| 1GB 서버에서 Kafka 3.5.1 싱글 노드 설치에 실패했습니다 — Java 8·ZooKeeper 레거시 기록 (0) | 2026.09.16 |
| [반자동 포스팅기] 1. 어디까지 자동으로 하고, 어디서 멈출 것인가 (0) | 2026.09.14 |
| 스팀잇 글쓰기 자동화 도구를 만들었습니다 — 같이 쓰실 분 찾습니다 (0) | 2026.08.26 |