ExecuteSQLRecord(JSON 출력) 결과를 Impala 테이블에 넣을 때, 레코드를 1000 건씩 잘라 INSERT INTO bronze.<table> (...) VALUES (...),(...) 문을 FlowFile 하나씩 만들어 PutSQL 로 보낸다. 수십만 건을 한 건씩 실행하는 PutSQL 보다 훨씬 빠르고, 테이블 이름은 FlowFile 속성 tablename 에서 받는다.
import org.apache.commons.io.IOUtils
import java.nio.charset.StandardCharsets
import groovy.json.JsonSlurper
def flowFile = session.get()
if (!flowFile) return
try {
def tableName = flowFile.getAttribute('tablename')
if (!tableName) {
log.error("tablename attribute가 없습니다.")
session.transfer(flowFile, REL_FAILURE)
return
}
def json = null
flowFile = session.write(flowFile, { inputStream, outputStream ->
json = IOUtils.toString(inputStream, StandardCharsets.UTF_8)
outputStream.write("".getBytes(StandardCharsets.UTF_8))
} as StreamCallback)
def records = new JsonSlurper().parseText(json)
def columns = ['epid', 'comp_cd', 'dept_cd', 'bt_type_cd', 'idm_del_yn', 'load_ts', 'part_ymd']
def esc = { v -> v?.toString()?.replace("'", "''") ?: '' }
records.collate(1000).eachWithIndex { batch, batchIdx ->
def newFlowFile = session.create(flowFile)
newFlowFile = session.write(newFlowFile, { inputStream, outputStream ->
def sb = new StringBuilder()
sb.append("INSERT INTO bronze.${tableName} (${columns.join(',')}) VALUES\n")
batch.eachWithIndex { r, idx ->
sb.append("(" + columns.collect { "'" + esc(r[it]) + "'" }.join(',') + ")")
sb.append(idx < batch.size() - 1 ? ",\n" : "\n")
}
outputStream.write(sb.toString().getBytes(StandardCharsets.UTF_8))
} as StreamCallback)
newFlowFile = session.putAllAttributes(newFlowFile, [
'batch.index': batchIdx.toString(),
'batch.count': batch.size().toString()
])
session.transfer(newFlowFile, REL_SUCCESS)
}
session.remove(flowFile)
} catch (Exception e) {
log.error("INSERT 생성 실패: ${e.message}", e)
session.transfer(flowFile, REL_FAILURE)
}
columns 목록만 테이블별로 바꾸면 된다. 숫자 컬럼(DECIMAL 등)은 따옴표 없이 넣도록 목록에 타입을 붙여 처리한다.; 나 ' 가 든 문자열은 SELECT 단계에서 replace() 로 거르거나 위 esc 로 이스케이프한다.INSERT OVERWRITE 를 배치 단위로 여러 번 실행하면 앞선 배치가 지워지므로, 배치 분할 방식에서는 INSERT INTO 를 쓰고 파티션 정리는 적재 전에 한 번만 한다.REL_FAILURE 로 보내고 LogAttribute 로 어느 배치가 실패했는지 남긴다.