SELECT
EMP_NUM,
NAME,
EMAIL,
-- 필요한 기타 컬럼들
FROM (
SELECT
EMP_NUM,
NAME,
EMAIL,
ORG_CD,
-- 동일 사번 내에서 본사(KR01)가 위로 오도록, 혹은 최신 등록일 순으로 정렬
ROW_NUMBER() OVER(PARTITION BY EMP_NUM ORDER BY (CASE WHEN ORG_CD = 'KR01' THEN 1 ELSE 2 END) ASC, UPDATE_DATE DESC) as RN
FROM
HR_TABLE
) T
WHERE
RN = 1; -- 1등 데이터만 추출 (중복 제거 완료)
import tkinter as tk
from tkinter import ttk, scrolledtext, messagebox
from confluent_kafka import Consumer, KafkaError
from confluent_kafka.admin import AdminClient
import threading
import json
import time
import os
import sys
import copy
# --- 설정 파일 경로 (절대 경로) ---
if getattr(sys, 'frozen', False):
BASE_DIR = os.path.dirname(sys.executable)
else:
BASE_DIR = os.path.dirname(os.path.abspath(__file__))
CONFIG_FILE = os.path.join(BASE_DIR, "kafka_profiles.json")
class KafkaConsumerApp:
def __init__(self, root):
self.root = root
self.root.title("Kafka Consumer UI (Auto-Save & Profile Modal)")
self.root.geometry("700x750")
# 종료 시 파일 저장
self.root.protocol("WM_DELETE_WINDOW", self.on_closing)
self.is_running = False
self.consumer_thread = None
# 데이터 저장소
self.profiles = {}
self.current_profile_name = None
# --- UI 변수 선언 (자동 저장을 위해 필수) ---
self.init_variables()
# --- UI 레이아웃 구성 ---
self.create_widgets()
# --- 초기 데이터 로드 ---
self.load_profiles_from_file()
def init_variables(self):
"""UI와 연동될 변수들을 선언하고, 변경 감지(Trace)를 설정합니다."""
# 각 변수가 변경될 때마다 self.on_input_change 함수가 실행되어 자동 저장됩니다.
self.var_profile_selection = tk.StringVar()
self.var_servers = tk.StringVar()
self.var_servers.trace_add("write", self.on_input_change)
self.var_use_security = tk.BooleanVar(value=False)
self.var_use_security.trace_add("write", self.on_input_change)
self.var_protocol = tk.StringVar(value="SASL_PLAINTEXT")
self.var_protocol.trace_add("write", self.on_input_change)
self.var_mech = tk.StringVar(value="PLAIN")
self.var_mech.trace_add("write", self.on_input_change)
self.var_user = tk.StringVar()
self.var_user.trace_add("write", self.on_input_change)
self.var_pw = tk.StringVar()
self.var_pw.trace_add("write", self.on_input_change)
self.var_topic = tk.StringVar()
self.var_topic.trace_add("write", self.on_input_change)
self.var_group = tk.StringVar()
self.var_group.trace_add("write", self.on_input_change)
self.var_offset = tk.StringVar(value="earliest")
self.var_offset.trace_add("write", self.on_input_change)
def create_widgets(self):
# 0. 프로필 관리 영역 (상단)
profile_frame = ttk.LabelFrame(self.root, text="Profile Management", padding="10")
profile_frame.pack(fill="x", padx=10, pady=5)
ttk.Label(profile_frame, text="Select Profile:").pack(side="left")
# 프로필 선택 콤보박스 (Readonly로 설정하여 오타 방지)
self.combo_profile = ttk.Combobox(profile_frame, textvariable=self.var_profile_selection, state="readonly", width=30)
self.combo_profile.pack(side="left", padx=5)
# 선택 이벤트 바인딩
self.combo_profile.bind("<<ComboboxSelected>>", self.on_profile_switched)
# 신규 생성 버튼
ttk.Button(profile_frame, text="➕ New Profile", command=self.open_new_profile_modal).pack(side="left", padx=5)
# 삭제 버튼
ttk.Button(profile_frame, text="🗑 Delete", command=self.delete_current_profile).pack(side="left", padx=5)
# 1. 연결 설정
conn_frame = ttk.LabelFrame(self.root, text="1. Connection & Security", padding="10")
conn_frame.pack(fill="x", padx=10, pady=5)
ttk.Label(conn_frame, text="Bootstrap Servers:").grid(row=0, column=0, sticky="w")
self.entry_servers = ttk.Entry(conn_frame, textvariable=self.var_servers, width=40)
self.entry_servers.grid(row=0, column=1, columnspan=3, padx=5, pady=2, sticky="w")
# Security Checkbox
self.chk_security = ttk.Checkbutton(
conn_frame,
text="Enable SASL/Security Protocol",
variable=self.var_use_security,
command=self.toggle_security_ui
)
self.chk_security.grid(row=1, column=0, columnspan=2, sticky="w", pady=(10, 5))
# Protocol & Mech
ttk.Label(conn_frame, text="Security Protocol:").grid(row=2, column=0, sticky="w")
self.combo_protocol = ttk.Combobox(conn_frame, textvariable=self.var_protocol, values=["SASL_PLAINTEXT", "SASL_SSL", "SSL"], state="readonly", width=18)
self.combo_protocol.grid(row=2, column=1, padx=5, pady=2, sticky="w")
ttk.Label(conn_frame, text="SASL Mechanism:").grid(row=2, column=2, sticky="w", padx=(10,0))
self.combo_mech = ttk.Combobox(conn_frame, textvariable=self.var_mech, values=["PLAIN", "SCRAM-SHA-256", "SCRAM-SHA-512"], state="readonly", width=18)
self.combo_mech.grid(row=2, column=3, padx=5, pady=2, sticky="w")
# User & PW
ttk.Label(conn_frame, text="Username:").grid(row=3, column=0, sticky="w")
self.entry_user = ttk.Entry(conn_frame, textvariable=self.var_user, width=21)
self.entry_user.grid(row=3, column=1, padx=5, pady=2, sticky="w")
ttk.Label(conn_frame, text="Password:").grid(row=3, column=2, sticky="w", padx=(10,0))
self.entry_pw = ttk.Entry(conn_frame, textvariable=self.var_pw, show="*", width=21)
self.entry_pw.grid(row=3, column=3, padx=5, pady=2, sticky="w")
# 2. 토픽 설정
topic_frame = ttk.LabelFrame(self.root, text="2. Topic & Consumer Group", padding="10")
topic_frame.pack(fill="x", padx=10, pady=5)
self.btn_load_topics = ttk.Button(topic_frame, text="🔄 Load Topics from Server", command=self.load_topics_from_server)
self.btn_load_topics.grid(row=0, column=0, columnspan=2, sticky="we", pady=(0, 5))
ttk.Label(topic_frame, text="Target Topic:").grid(row=1, column=0, sticky="w")
self.combo_topic = ttk.Combobox(topic_frame, textvariable=self.var_topic, width=40)
self.combo_topic.grid(row=1, column=1, padx=5, pady=2)
self.lbl_topic_count = ttk.Label(topic_frame, text="(Total: -)", foreground="gray")
self.lbl_topic_count.grid(row=1, column=2, padx=5, sticky="w")
ttk.Label(topic_frame, text="Group ID:").grid(row=2, column=0, sticky="w")
self.entry_group = ttk.Entry(topic_frame, textvariable=self.var_group, width=42)
self.entry_group.grid(row=2, column=1, padx=5, pady=2)
ttk.Label(topic_frame, text="Offset Reset:").grid(row=3, column=0, sticky="w")
self.combo_offset = ttk.Combobox(topic_frame, textvariable=self.var_offset, values=["earliest", "latest"], state="readonly", width=40)
self.combo_offset.grid(row=3, column=1, padx=5, pady=2)
# 3. 제어 버튼
btn_frame = ttk.Frame(self.root, padding="5")
btn_frame.pack(fill="x", padx=10)
self.btn_start = ttk.Button(btn_frame, text="▶ Start Consuming", command=self.start_consuming)
self.btn_start.pack(side="left", padx=5, fill="x", expand=True)
self.btn_stop = ttk.Button(btn_frame, text="⏹ Stop", command=self.stop_consuming, state="disabled")
self.btn_stop.pack(side="left", padx=5, fill="x", expand=True)
self.btn_clear = ttk.Button(btn_frame, text="🧹 Clear Log", command=self.clear_log)
self.btn_clear.pack(side="right", padx=5)
# 4. 로그
log_frame = ttk.LabelFrame(self.root, text="Real-time Messages", padding="5")
log_frame.pack(fill="both", expand=True, padx=10, pady=5)
self.txt_log = scrolledtext.ScrolledText(log_frame, state="disabled", height=10)
self.txt_log.pack(fill="both", expand=True)
# ================= 프로필 관리 로직 =================
def load_profiles_from_file(self):
"""파일에서 프로필 로드"""
if os.path.exists(CONFIG_FILE):
try:
with open(CONFIG_FILE, "r", encoding="utf-8") as f:
data = json.load(f)
self.profiles = data.get("profiles", {})
last_used = data.get("last_used", "")
self.update_combo_list()
# 마지막 사용 프로필 복원
if last_used and last_used in self.profiles:
self.combo_profile.set(last_used)
self.on_profile_switched(None)
elif self.profiles:
# 없으면 첫번째꺼
first = list(self.profiles.keys())[0]
self.combo_profile.set(first)
self.on_profile_switched(None)
else:
# 프로필이 아예 없으면 기본값 생성 유도
self.open_new_profile_modal(force_first=True)
except Exception as e:
print(f"로드 실패: {e}")
else:
# 파일이 없으면 바로 모달 띄우기
self.root.after(100, lambda: self.open_new_profile_modal(force_first=True))
def save_profiles_to_file(self):
"""현재 프로필 목록과 상태를 파일에 저장 (프로그램 종료나 프로필 변경 시 호출)"""
data = {
"last_used": self.current_profile_name if self.current_profile_name else "",
"profiles": self.profiles
}
try:
with open(CONFIG_FILE, "w", encoding="utf-8") as f:
json.dump(data, f, indent=4, ensure_ascii=False)
print("파일 저장 완료")
except Exception as e:
messagebox.showerror("Error", f"저장 실패: {e}")
def on_input_change(self, *args):
"""
[자동 저장 핵심]
UI 입력값이 변경될 때마다 호출되어 self.profiles 딕셔너리를 즉시 업데이트합니다.
파일 쓰기(I/O)는 종료/스위칭 때만 발생시켜 성능을 확보합니다.
"""
if not self.current_profile_name:
return
# 현재 UI 값들을 딕셔너리로 만듦
updated_data = {
"bootstrap_servers": self.var_servers.get(),
"use_security": self.var_use_security.get(),
"security_protocol": self.var_protocol.get(),
"sasl_mechanism": self.var_mech.get(),
"sasl_username": self.var_user.get(),
"sasl_password": self.var_pw.get(),
"topic": self.var_topic.get(),
"group_id": self.var_group.get(),
"auto_offset_reset": self.var_offset.get()
}
# 메모리 상의 프로필 정보 업데이트
self.profiles[self.current_profile_name] = updated_data
# 보안 UI 상태도 실시간 반영
self.toggle_security_ui()
def on_profile_switched(self, event):
"""콤보박스에서 프로필 선택 시 실행"""
new_profile_name = self.combo_profile.get()
if not new_profile_name or new_profile_name not in self.profiles:
return
# 이전 프로필 파일 저장 (안전장치)
if self.current_profile_name:
self.save_profiles_to_file()
self.current_profile_name = new_profile_name
data = self.profiles[new_profile_name]
# UI 변수 업데이트 (Trace가 발동되지 않도록 값을 직접 세팅하는 것이 좋으나,
# 여기선 Trace가 발동되어도 같은 값이므로 상관없음. 단, 무한루프 주의)
# 값 적용
self.var_servers.set(data.get('bootstrap_servers', 'localhost:9092'))
self.var_use_security.set(data.get('use_security', False))
self.var_protocol.set(data.get('security_protocol', 'SASL_PLAINTEXT'))
self.var_mech.set(data.get('sasl_mechanism', 'PLAIN'))
self.var_user.set(data.get('sasl_username', ''))
self.var_pw.set(data.get('sasl_password', ''))
self.var_topic.set(data.get('topic', ''))
self.var_group.set(data.get('group_id', ''))
self.var_offset.set(data.get('auto_offset_reset', 'earliest'))
# UI 상태 갱신
self.toggle_security_ui()
print(f"프로필 로드 완료: {new_profile_name}")
def update_combo_list(self):
keys = sorted(list(self.profiles.keys()))
self.combo_profile['values'] = keys
# ================= 모달 창 (팝업) 로직 =================
def open_new_profile_modal(self, force_first=False):
"""신규 프로필 생성 모달 창"""
modal = tk.Toplevel(self.root)
modal.title("New Profile")
modal.geometry("400x250")
modal.transient(self.root) # 메인 창 위에 뜨게 함
modal.grab_set() # 모달 (다른 창 클릭 불가)
# 내용
ttk.Label(modal, text="New Profile Name:").pack(pady=(20, 5))
entry_name = ttk.Entry(modal, width=30)
entry_name.pack(pady=5)
entry_name.focus()
# 복사 기능
ttk.Label(modal, text="Copy Settings From (Optional):").pack(pady=(15, 5))
existing_profiles = sorted(list(self.profiles.keys()))
combo_copy = ttk.Combobox(modal, values=existing_profiles, state="readonly", width=28)
combo_copy.pack(pady=5)
# 현재 선택된 프로필이 있다면 그걸 기본값으로 선택
if self.current_profile_name in existing_profiles:
combo_copy.set(self.current_profile_name)
def create_action():
name = entry_name.get().strip()
if not name:
messagebox.showwarning("Warning", "프로필 이름을 입력해주세요.", parent=modal)
return
if name in self.profiles:
messagebox.showerror("Error", "이미 존재하는 이름입니다.", parent=modal)
return
# 데이터 생성 로직
new_data = {}
copy_source = combo_copy.get()
if copy_source and copy_source in self.profiles:
# [복사 기능] 깊은 복사로 기존 설정 가져오기
new_data = copy.deepcopy(self.profiles[copy_source])
else:
# 기본값
new_data = {
"bootstrap_servers": "localhost:9092",
"use_security": False,
"group_id": "my-group",
"auto_offset_reset": "earliest"
}
# 프로필 등록
self.profiles[name] = new_data
self.current_profile_name = name
# 콤보박스 갱신 및 선택
self.update_combo_list()
self.combo_profile.set(name)
# 실제 값 로드
self.on_profile_switched(None)
# 파일 저장
self.save_profiles_to_file()
modal.destroy()
btn_text = "Create & Start" if force_first else "Create"
ttk.Button(modal, text=btn_text, command=create_action).pack(pady=20)
def delete_current_profile(self):
"""현재 프로필 삭제"""
name = self.current_profile_name
if not name:
return
if messagebox.askyesno("Delete", f"프로필 '{name}'을(를) 정말 삭제하시겠습니까?"):
del self.profiles[name]
self.current_profile_name = None
self.combo_profile.set("")
# 값 초기화
self.var_servers.set("")
self.var_topic.set("")
self.update_combo_list()
self.save_profiles_to_file()
# 남은 프로필이 있으면 첫번째꺼 로드
if self.profiles:
first = list(self.profiles.keys())[0]
self.combo_profile.set(first)
self.on_profile_switched(None)
# ================= 기타 UI 로직 =================
def toggle_security_ui(self, *args):
is_enabled = self.var_use_security.get()
state = "readonly" if is_enabled else "disabled"
entry_state = "normal" if is_enabled else "disabled"
self.combo_protocol.config(state=state)
self.combo_mech.config(state=state)
self.entry_user.config(state=entry_state)
self.entry_pw.config(state=entry_state)
def get_kafka_config(self):
conf = {'bootstrap.servers': self.var_servers.get()}
if self.var_use_security.get():
conf.update({
'security.protocol': self.var_protocol.get(),
'sasl.mechanism': self.var_mech.get(),
'sasl.username': self.var_user.get(),
'sasl.password': self.var_pw.get()
})
return conf
def load_topics_from_server(self):
conf = self.get_kafka_config()
if not conf.get('bootstrap.servers'):
messagebox.showerror("Error", "Bootstrap Servers 주소를 입력하세요.")
return
self.btn_load_topics.config(state="disabled", text="Connecting...")
self.root.update()
try:
admin_client = AdminClient(conf)
meta = admin_client.list_topics(timeout=10)
topics = list(meta.topics.keys())
topics.sort()
if topics:
self.combo_topic['values'] = topics
current = self.var_topic.get()
if not current or current not in topics:
self.combo_topic.current(0)
self.lbl_topic_count.config(text=f"(Total: {len(topics)})")
messagebox.showinfo("Success", f"{len(topics)}개의 토픽 로드 완료")
else:
self.lbl_topic_count.config(text="(Total: 0)")
messagebox.showwarning("Warning", "토픽이 없습니다.")
except Exception as e:
messagebox.showerror("Error", f"실패: {e}")
finally:
self.btn_load_topics.config(state="normal", text="🔄 Load Topics from Server")
def on_closing(self):
if self.is_running:
if messagebox.askokcancel("Quit", "컨슈머 실행 중입니다. 종료할까요?"):
self.is_running = False
self.save_profiles_to_file()
self.root.destroy()
else:
self.save_profiles_to_file()
self.root.destroy()
# ================= 컨슈밍 로직 (스레드) =================
def start_consuming(self):
conf = self.get_kafka_config()
topic = self.var_topic.get()
group_id = self.var_group.get()
offset_reset = self.var_offset.get()
if not topic or not group_id:
messagebox.showerror("Error", "토픽과 그룹 ID를 확인해주세요.")
return
conf.update({
'group.id': group_id,
'auto.offset.reset': offset_reset,
'enable.auto.commit': True
})
self.is_running = True
self.toggle_ui_state(running=True)
self.consumer_thread = threading.Thread(
target=self.consume_loop, args=(conf, topic), daemon=True
)
self.consumer_thread.start()
def stop_consuming(self):
self.is_running = False
self.log_message("\n[System] Stopping...")
def consume_loop(self, conf, topic):
try:
consumer = Consumer(conf)
consumer.subscribe([topic])
self.log_message(f"[System] Connected to {conf['bootstrap.servers']}")
while self.is_running:
msg = consumer.poll(1.0)
if msg is None: continue
if msg.error():
if msg.error().code() != KafkaError._PARTITION_EOF:
self.log_message(f"[Error] {msg.error()}")
continue
try:
val = msg.value().decode('utf-8')
try:
js = json.loads(val)
val = json.dumps(js, indent=2, ensure_ascii=False)
except: pass
log = f"--- [Offset: {msg.offset()}] ---\n{val}"
self.root.after(0, self.log_message, log)
except Exception as e:
self.root.after(0, self.log_message, f"[Decode Error] {e}")
except Exception as e:
self.root.after(0, self.log_message, f"[Fatal Error] {e}")
finally:
if 'consumer' in locals(): consumer.close()
self.root.after(0, self.toggle_ui_state, False)
self.root.after(0, self.log_message, "[System] Stopped.")
def log_message(self, msg):
self.txt_log.config(state="normal")
self.txt_log.insert(tk.END, msg + "\n")
self.txt_log.see(tk.END)
self.txt_log.config(state="disabled")
def clear_log(self):
self.txt_log.config(state="normal")
self.txt_log.delete(1.0, tk.END)
self.txt_log.config(state="disabled")
def toggle_ui_state(self, running):
state = "disabled" if running else "normal"
# 실행 중엔 프로필 변경 불가
self.combo_profile.config(state=state if not running else "disabled") # Readonly handling
self.entry_servers.config(state=state)
self.btn_start.config(state="disabled" if running else "normal")
self.btn_stop.config(state="normal" if running else "disabled")
if __name__ == "__main__":
root = tk.Tk()
app = KafkaConsumerApp(root)
root.mainloop()
spring:
kafka:
properties:
# >- 는 "줄바꿈을 공백 하나로 치환하고 끝 공백 제거"라는 뜻입니다.
sasl.jaas.config: >-
org.apache.kafka.common.security.scram.ScramLoginModule required
username="myuser"
password="mypassword";
@Configuration
public class KafkaConfig {
@Bean
public Map<String, Object> producerConfigs() {
Map<String, Object> props = new HashMap<>();
// ... 다른 설정들 ...
props.put("security.protocol", "SASL_PLAINTEXT");
props.put("sasl.mechanism", "SCRAM-SHA-256");
// Java 코드에서 문자열로 넣으면 yml 파싱 에러 걱정이 없습니다.
String jaasConfig = String.format(
"org.apache.kafka.common.security.scram.ScramLoginModule required username=\"%s\" password=\"%s\";",
"myuser", "mypassword"
);
props.put("sasl.jaas.config", jaasConfig);
return props;
}
}
