JDBC INSERT 로 대량 적재하면 느리다. PostgreSQL 은 COPY 명령으로 스트림을 통째로 받아들이며, JDBC 드라이버는 이를 CopyManager 로 노출한다. Spark 파티션마다 커넥션을 열고 파티션 데이터를 스트림으로 밀어 넣는 방식이 실무에서 널리 쓰인다.
파티션 단위로 처리하는 것이 요점이다. 행 하나마다 커넥션을 여는 코드는 어떤 튜닝으로도 살릴 수 없다.
df.foreachPartition((ForeachPartitionFunction<Row>) rows -> {
try (Connection conn = DriverManager.getConnection(url, props)) {
CopyManager cm = conn.unwrap(BaseConnection.class).getCopyAPI();
PipedOutputStream out = new PipedOutputStream();
PipedInputStream in = new PipedInputStream(out, 1 << 16);
Thread writer = new Thread(() -> {
try (OutputStreamWriter w = new OutputStreamWriter(out, StandardCharsets.UTF_8)) {
while (rows.hasNext()) {
Row r = rows.next();
w.write(toDelimitedLine(r, "\t"));
w.write('\n');
}
} catch (IOException e) {
throw new RuntimeException(e);
}
});
writer.start();
cm.copyIn("COPY my_table FROM STDIN WITH (FORMAT text, DELIMITER E'\\t')", in);
writer.join();
}
});
conn.unwrap(BaseConnection.class) 를 쓴다. 커넥션 풀을 거치면 실제 객체가 래퍼이므로 (BaseConnection) conn 형변환은 실패한다.
구분자와 개행이 값 안에 들어 있으면 그대로 깨진다. 텍스트 포맷이라면 이스케이프 규칙을 지키거나, 아예 CSV 포맷을 쓰고 인용 처리를 맡긴다.
COPY my_table FROM STDIN WITH (FORMAT csv, DELIMITER ',', QUOTE '"', NULL '');
NULL 표현도 정해야 한다. 텍스트 포맷의 기본 NULL 표기는 \N 이며, 빈 문자열과 구분된다. 원본이 빈 문자열과 NULL 을 모두 갖는다면 어느 쪽으로 넣을지 결정하고 명시한다.
인코딩은 UTF-8 로 맞춘다. 세션 인코딩과 다르면 한글이 깨진다.
COPY 는 한 번의 트랜잭션으로 처리되므로 중간에 한 행이라도 타입이 맞지 않으면 그 파티션 전체가 실패한다. 실무에서는 다음 중 하나를 택한다. 스테이징 테이블에 모두 문자열로 받아 적재 후 검증하고 변환하거나, 사전에 스키마를 검증해 불량 행을 걸러낸다.
재실행을 고려한다면 파티션 단위로 적재 상태를 기록하거나, 스테이징 테이블에 넣고 마지막에 한 번에 교체하는 방식을 쓴다.
BEGIN;
TRUNCATE my_table_stage;
-- COPY ...
INSERT INTO my_table SELECT * FROM my_table_stage;
COMMIT;
파티션 수만큼 동시에 커넥션이 열린다. 파티션이 200개면 PostgreSQL 에 동시 접속 200개가 들어온다. max_connections 와 서버 자원을 넘지 않도록 적재 직전에 파티션 수를 줄인다.
df.coalesce(16).foreachPartition(...);
기존 Scala 로 작성된 유틸리티가 있다면 Java 에서 싱글턴 객체를 통해 호출한다. Scala 의 object 는 컴파일 후 클래스명$.MODULE$ 로 접근한다.
InputStream in = PostgresUtil$.MODULE$.rowsToInputStream(rows, "\t");
다만 Scala 컬렉션과 Java 컬렉션은 타입이 다르므로 변환이 필요하고, Scala 버전이 바뀌면 시그니처가 달라진다. 새로 만든다면 Java 쪽에서 직접 구현하는 편이 유지보수가 쉽다.
Spark 2.4.x 는 Scala 2.11 과 2.12 빌드가 있고, Spark 3.x 는 2.12 와 2.13, Spark 4.0 은 2.13 만 제공한다. 프로젝트의 scala-library, spark-core_<scala>, 테스트 라이브러리의 Scala 접미사를 모두 같은 계열로 맞춘다. requires scala version: 2.11.7 같은 빌드 오류는 라이브러리가 기대하는 Scala 패치 버전과 프로젝트의 버전이 어긋났다는 뜻이며, 대개 Scala 버전을 올리는 쪽으로 정리한다.