ListFile → FetchFile → RouteOnAttribute(크기 0 제외) → SplitText → ExtractText → UpdateAttribute → AttributesToJSON → PutKudu 로 NAS 로그 파일을 적재하면서, DB · NAS · Kafka 세 원천의 적재 이력을 한 곳에서 관리하고 실패 건만 골라 재수행할 수 있게 한다. 대상 테이블은 Upsert 로 적재하므로 중복 적재는 이미 막혀 있다.
fragment.index 는 한 배치 안에서만 유일하다. 같은 테이블에 배치가 여러 번 돌면 0 부터 다시 시작해 PK 가 충돌하므로, 배치를 구분하는 값이 키에 있어야 한다.MAX(seq)+1 방식은 NiFi 가 FlowFile 을 동시 처리할 때 충돌하고 매 건 SELECT 가 든다. 대신 배치 시작 시각(base_dts)을 키에 넣고 fragment.index 를 그대로 detail_seq 로 쓴다.src_type 으로 구분하고, detail 에 원천별 컬럼을 nullable 로 둔다.CREATE TABLE load_log_master (
src_type STRING, -- DB / NAS / KAFKA
regn STRING,
db_name STRING,
table_name STRING,
base_dts STRING, -- 배치 시작 시각 yyyyMMddHHmmss
base_dt STRING, -- 기준일 yyyyMMdd, 조회용
status STRING, -- RUNNING / SUCCESS / PARTIAL / FAIL
src_count INT,
load_count INT,
error_count INT,
start_ts TIMESTAMP,
end_ts TIMESTAMP,
created_dt TIMESTAMP,
updated_dt TIMESTAMP,
PRIMARY KEY (src_type, regn, db_name, table_name, base_dts)
)
PARTITION BY HASH(src_type, regn) PARTITIONS 4
STORED AS KUDU;
CREATE TABLE load_log_detail (
src_type STRING,
regn STRING,
db_name STRING,
table_name STRING,
base_dts STRING,
detail_seq INT, -- SplitText fragment.index
base_dt STRING,
status STRING,
error_message STRING,
error_processor STRING,
retry_count INT,
file_path STRING, -- NAS
file_name STRING,
line_content STRING,
kafka_topic STRING, -- KAFKA
kafka_category STRING,
kafka_svc STRING,
kafka_offset BIGINT,
kafka_partition INT,
record_key STRING, -- DB
created_dt TIMESTAMP,
updated_dt TIMESTAMP,
PRIMARY KEY (src_type, regn, db_name, table_name, base_dts, detail_seq)
)
PARTITION BY HASH(src_type, regn) PARTITIONS 4
STORED AS KUDU;
base_dts 는 NiFi 에서 ${now():format('yyyyMMddHHmmss')} 로 만들어 배치 안의 모든 FlowFile 에 UpdateAttribute 로 붙인다.
-- 기준일 전체 현황
SELECT * FROM load_log_master WHERE base_dt = '20260225';
-- 특정 배치의 실패 건
SELECT * FROM load_log_detail
WHERE src_type = 'KAFKA' AND regn = 'R01' AND db_name = 'event_db'
AND table_name = 'click' AND base_dts = '20260225143000' AND status = 'FAIL';
-- 재수행 대상
SELECT * FROM load_log_detail WHERE status = 'FAIL' AND retry_count < 3;
AttributesToJSON 이후에 실패한 건은 detail 의 line_content 로 바로 재처리할 수 있지만, 그 이전 단계(FetchFile · SplitText · ExtractText)에서 실패하면 가공된 내용이 없어 원본 파일을 다시 읽어야 한다. 원본이 계속 append 되는 파일이면 ListFile + FetchFile 은 매번 전체를 읽으므로 라인 내용(line_content)을 detail 에 보존해 두어야 실패 건만 다시 넣을 수 있다.