수십만 건을 PutSQL 로 INSERT 하면 오래 걸린다. 대신 ExecuteSQLRecord → UpdateRecord(load_ts, part_ymd 추가) → ConvertRecord(Parquet) → PutHDFS 로 파일을 올리고, Impala 로 파티션을 등록한다. 대상은 bronze 외부 테이블(/warehouse/tablespace/external/hive/bronze.db/<table>/part_ymd=<yyyyMMdd>)이다.
PutHDFS 가 /tmp/nifi/bronze.db/${tablename}/${part_ymd} 에 쓰고, ExecuteScript(Groovy) 가 Impala JDBC 로 ADD PARTITION → LOAD DATA INPATH ... OVERWRITE → REFRESH 를 실행한다. OVERWRITE 가 기존 파티션 파일을 지우므로 재처리 때 파일이 쌓이지 않는다.
import java.sql.Connection
import java.sql.Statement
import org.apache.nifi.dbcp.DBCPService
def flowFile = session.get()
if (!flowFile) return
def tableName = flowFile.getAttribute('tablename')
def partYmd = flowFile.getAttribute('part_ymd')
def tmpPath = flowFile.getAttribute('tmp.path')?.replaceAll('/+$', '')
def database = flowFile.getAttribute('database')
def poolUuid = flowFile.getAttribute('impala_connection_uuid')?.trim()
Connection conn = null; Statement stmt = null
try {
def dbcp = context.controllerServiceLookup.getControllerService(poolUuid) as DBCPService
conn = dbcp.getConnection(); stmt = conn.createStatement()
stmt.execute("ALTER TABLE ${database}.${tableName} ADD IF NOT EXISTS PARTITION (part_ymd='${partYmd}')")
stmt.execute("LOAD DATA INPATH '${tmpPath}/${partYmd}' OVERWRITE INTO TABLE ${database}.${tableName} PARTITION (part_ymd='${partYmd}')")
stmt.execute("REFRESH ${database}.${tableName}")
session.transfer(flowFile, REL_SUCCESS)
} catch (Exception e) {
log.error("실행 실패 [${database}.${tableName} / ${partYmd}]: ${e.message}", e)
session.transfer(flowFile, REL_FAILURE)
} finally {
try { stmt?.close() } catch (Exception ignored) {}
try { conn?.close() } catch (Exception ignored) {}
}
database 는 UpdateAttribute 에서 #{etl_db} 로, impala_connection_uuid 는 ImpalaConnectionPool 서비스의 ID 로 넣는다. Statement 는 한 번 만들어 재사용한다.
LOAD DATA 는 HDFS 안에서 rename 이므로 임시 경로와 웨어하우스가 같은 nameservice 여야 한다. Federation 으로 nameservice 가 둘 이상이면 실패하므로 NiFi 가 옮긴다.
PutHDFS(/tmp/...) → MoveHDFS → ExecuteScript(ADD PARTITION + REFRESH)
| MoveHDFS 속성 | 값 |
|---|---|
| Hadoop Configuration Resources | /etc/hadoop/conf/core-site.xml,/etc/hadoop/conf/hdfs-site.xml |
| Kerberos Principal / Keytab | ETL 계정 |
| Input Directory | ${tmp.path}/${part_ymd} |
| Output Directory | ${hdfs.path}/part_ymd=${part_ymd} |
| Operation | move |
| Conflict Resolution | replace |
MoveHDFS 는 기존 파일을 지우지 않고 파일을 추가하므로, 파티션 전체를 덮어쓰려면 앞에 DeleteHDFS 로 part_ymd= 디렉터리를 비우거나 PutHDFS 를 처음부터 최종 경로에 replace 로 쓰고 LOAD DATA 를 생략한다. 개발(단일 클러스터)과 운영(다중 nameservice)이 다르면 운영 기준으로 맞춘다.
AccessControlException 이 나면 NiFi 실행 계정의 keytab 이 아니라 프로세서에 지정한 principal 로 /tmp/nifi 와 웨어하우스 경로에 쓰기 권한이 있는지 본다. hdfs dfs -stat %F <path> · hdfs dfs -ls 를 그 keytab 으로 kinit 한 뒤 확인한다.
load_ts · part_ymd 추가하는 설정.