로그 파일을 NiFi 로 수집해 Kudu 에 넣을 때, 같은 파일이 다시 들어오거나 재처리를 돌리면 같은 행이 두 번 적재된다. Kudu 는 기본 키가 같으면 덮어쓰므로 기본 키만 잘 잡으면 재적재가 멱등해진다. 그런데 로그의 자연 키(시각 · 사용자 · 서비스 정도)만으로는 한 건을 특정하지 못하는 경우가 있다. 같은 시각에 같은 사용자가 낸 서로 다른 호출이 구분되지 않는다.
이때 나머지 컬럼 전체를 재료로 삼아 만든 내용 해시를 기본 키의 마지막 컬럼으로 넣는다. 내용이 같으면 해시가 같아 덮어쓰기가 되고, 내용이 한 글자라도 다르면 다른 행으로 남는다. 별도의 중복 판정 쿼리나 staging 테이블이 필요 없다.
기본 키 앞쪽에는 조회에 쓰는 자연 키를 두고, 마지막에 row_hash 를 붙인다. 해시 파티션도 기본 키 전체로 잡아 쓰기가 한쪽으로 몰리지 않게 한다.
CREATE TABLE bronze.nas_copilot_call (
ts STRING NOT NULL,
epid STRING NOT NULL,
serviceid STRING NOT NULL,
functiontype STRING NOT NULL,
clickid STRING NOT NULL,
row_hash STRING NOT NULL,
response STRING NULL,
userid STRING NULL,
comp STRING NULL,
suborg STRING NULL,
llmurl STRING NULL,
llmtype STRING NULL,
input STRING NULL,
output STRING NULL,
func_name STRING NULL,
func_name_en STRING NULL,
analyticsyn STRING NULL,
devicetype STRING NULL,
PRIMARY KEY (ts, epid, serviceid, functiontype, clickid, row_hash)
)
PARTITION BY HASH (ts, epid, serviceid, functiontype, clickid, row_hash) PARTITIONS 6
STORED AS KUDU;
ExecuteScript 에 Groovy 로 넣는다. AttributesToJSON 앞에 두어 row_hash 속성이 JSON 목록에 포함되게 한다.
import java.security.MessageDigest
def flowFile = session.get()
if (flowFile == null) return
// row_hash 를 제외한 컬럼 전부. 알파벳 순으로 고정한다.
def cols = ['analyticsyn','clickid','comp','devicetype','epid',
'func_name','func_name_en','functiontype','input',
'llmtype','llmurl','output','response',
'serviceid','suborg','ts','userid']
def concatStr = cols.collect { c ->
def v = flowFile.getAttribute(c)
// null 과 빈 문자열을 모두 \N 으로 통일한다.
return (v == null || v.isEmpty()) ? '\\N' : v
}.join('|')
def hash = MessageDigest.getInstance('MD5')
.digest(concatStr.getBytes('UTF-8'))
.encodeHex()
.toString()
flowFile = session.putAttribute(flowFile, 'row_hash', hash)
session.transfer(flowFile, REL_SUCCESS)
해시는 같은 입력에 대해 언제 어디서 계산해도 같은 값이 나와야 한다. 그렇지 않으면 재적재가 멱등해지기는커녕 매번 새 행을 만든다. 규칙을 네 가지로 못 박아 두고 문서에 남긴다.
null 과 빈 문자열을 구분하지 않고 둘 다 \N 으로 쓴다. 파서가 없는 필드를 null 로 줄 때와 빈 문자열로 줄 때가 섞이기 때문이다.| 로 잇는다. 값 안에 나타나지 않는 문자를 쓴다. 값에 섞일 수 있는 문자를 쓰면 a|b + c 와 a + b|c 가 같은 해시가 된다.getBytes('UTF-8') 로 명시한다. JVM 기본 인코딩에 맡기면 노드마다 값이 달라진다.이미 들어간 데이터를 Impala 나 Trino 에서 다시 검산하거나, 백필을 SQL 로 돌릴 때가 있다. 이때 SQL 쪽 표현이 Groovy 쪽과 한 글자라도 다르면 같은 행이 다른 해시를 받는다. 빈 값 정규화는 아래 표현과 짝을 이룬다.
coalesce(nullif(col, ''), '\\N')
SQL 쪽에서 통째로 계산하려면 같은 순서 · 같은 구분자로 이어 붙인 뒤 해시를 건다.
SELECT lower(hex(md5(concat_ws('|',
coalesce(nullif(analyticsyn, ''), '\\N'),
coalesce(nullif(clickid, ''), '\\N'),
coalesce(nullif(comp, ''), '\\N')
-- 나머지 컬럼도 같은 순서로
)))) AS row_hash
FROM bronze.nas_copilot_call;
AttributesToJSON 목록에 반드시 넣는다. 속성으로만 만들어 두고 JSON 목록에서 빠뜨리면 PutKudu 가 기본 키 없이 적재를 시도해 실패한다.