Engineering Note

Connection은 하나, Table은 스레드마다: Apache Fluss Java Client 실습

Flink Quickstart의 Primary Key Table과 Log Table을 Apache Fluss 0.9.1 Java Client로 읽고 쓰며 객체 수명 주기를 정리합니다.

2026년 8월 1일 · Pletor Engineering apache-flussjavaquickstartstreaming
안정적인 중심에서 여러 경로가 부드럽게 이어지는 종이 콜라주 일러스트
하나의 안정적인 연결에서 여러 작업을 차분히 시작하는 모습을 표현한 일러스트

Flink SQL로 테이블을 만들었다면, 다음 질문은 애플리케이션 코드에서 그 테이블을 어떻게 안전하게 다룰 것인가입니다. Apache Fluss Java Client의 기준은 간단합니다. 애플리케이션마다 Connection은 하나를 공유하고, 작업을 시작하는 스레드마다 Table 또는 Admin을 새로 만듭니다.

이 글은 Flink Quickstart에서 만든 demo.customer_profiledemo.order_events를 그대로 사용합니다. 전자는 같은 키의 최신 상태를 보관하는 Primary Key Table이고, 후자는 이벤트가 계속 쌓이는 Log Table입니다. Java Client도 이 두 모델을 하나의 쓰기 API로 섞지 않습니다.

Konduo는 Fluss 같은 데이터 인프라를 관리·운영·모니터링하는 도구입니다. Java Client의 객체 수명 주기를 이해해 두면, 애플리케이션 코드와 운영 도구가 같은 클러스터에 연결할 때 연결을 과도하게 만들거나 공유하면 안 되는 객체를 재사용하는 실수를 줄일 수 있습니다.

이 글에서 확인할 것

  • Connection은 공유하지만 TableAdmin은 공유하지 않는 이유
  • GenericRow로 Primary Key Table을 upsert하고 키로 현재 상태를 찾는 방법
  • Log Table에 append한 뒤 LogScanner로 버킷을 구독하는 방법
  • 행 형식 API와 POJO 기반 Typed API를 고르는 기준

시작하기 전에: Quickstart 테이블을 준비한다

먼저 Apache Fluss Flink Quickstart를 끝까지 실행해 demo 데이터베이스와 다음 테이블을 만듭니다. 이 글은 테이블을 새로 만들거나 기존 데이터를 지우지 않습니다.

demo.customer_profile  Primary Key Table
  customer_id     INT PRIMARY KEY
  name            STRING
  membership      STRING
  account_balance DECIMAL(15, 2)

demo.order_events      Log Table
  order_id             BIGINT
  customer_id          INT
  total_price          DECIMAL(15, 2)
  ordered_on           DATE
  order_priority       STRING
  clerk                STRING

fluss-client 0.9.1-incubating을 프로젝트에 추가합니다. 이 글의 API와 동작은 Apache Fluss 0.9.1-incubating 기준입니다.

<dependency>
  <groupId>org.apache.fluss</groupId>
  <artifactId>fluss-client</artifactId>
  <version>0.9.1-incubating</version>
</dependency>

본문의 코드 조각과 별도로, 바로 실행할 수 있는 두 Java 파일과 Maven 설정 파일도 제공합니다. 세 파일을 같은 디렉터리에 내려받습니다.

# JDK 17 이상과 Maven이 필요하다.
# 세 파일을 같은 빈 디렉터리에 둔 뒤 각각 실행한다.
MAVEN_OPTS='--add-opens=java.base/java.nio=ALL-UNNAMED' \
  mvn compile exec:java -Dexec.mainClass=PrimaryKeyTableExample

MAVEN_OPTS='--add-opens=java.base/java.nio=ALL-UNNAMED' \
  mvn compile exec:java -Dexec.mainClass=LogTableExample

JDK 17 이상이면 Arrow 접근을 허용한다

fluss-client의 쓰기 경로는 Arrow의 메모리 기능을 사용합니다. JDK 17 이상에서는 java.nio 접근을 열지 않으면 첫 append 또는 upsert에서 초기화 오류가 날 수 있습니다. 애플리케이션을 시작하는 JVM 옵션에 다음을 추가합니다.

--add-opens=java.base/java.nio=ALL-UNNAMED

Connection은 공유하고 Table은 작업마다 만든다

Connection은 Coordinator·TabletServer 연결과 연결 설정의 시작점입니다. 스레드 안전하므로 애플리케이션에서 하나만 만들어 공유합니다. 반면 TableAdmin은 스레드 안전하지 않습니다. 필요한 스레드에서 만들고, 작업이 끝나면 닫습니다. 캐시나 풀로 재사용하지 않는 편이 안전합니다.

Configuration config = new Configuration();
config.setString("bootstrap.servers", "localhost:9123");

// 애플리케이션을 시작할 때 한 번 만든다.
Connection connection = ConnectionFactory.createConnection(config);

// 이 Table 인스턴스는 이 작업을 실행하는 스레드에서만 사용한다.
try (Table profiles = connection.getTable(
        TablePath.of("demo", "customer_profile"))) {
    System.out.println(profiles.getTableInfo());
}

// 애플리케이션 종료 시점에 connection.close()를 호출한다.

이후 예제는 위에서 만든 connection이 애플리케이션 종료 전까지 살아 있다고 가정합니다. Admin도 같은 원칙을 따릅니다. connection.getAdmin()으로 얻되, 데이터베이스·테이블 관리 작업을 수행하는 스레드 안에서 만들고 오래 보관하지 않습니다. Admin 메서드는 CompletableFuture를 돌려주므로, 짧은 관리 작업은 get()으로 기다릴 수 있고 비동기 흐름에서는 thenApply·exceptionally로 연결할 수 있습니다.

하나의 Connection이 스레드별 Table과 Admin으로 이어지고 Table에서 writer, lookuper, scanner가 만들어지는 구조도
공유 범위는 Connection까지입니다. Table·Admin과 그 아래의 작업 객체는 필요한 작업과 스레드의 범위에서 만듭니다.

GenericRow로 Primary Key Table의 최신 상태를 갱신하고 확인한다

GenericRow는 필드 순서로 값을 넣는 저수준 행 형식입니다. 테이블 스키마를 직접 알고 있고 변환 비용을 줄이고 싶을 때 적합합니다. 다만 STRINGBinaryString, DECIMALDecimal처럼 Fluss 내부 형식을 넣어야 합니다. 보통의 Java 객체를 그대로 넣는 API가 아닙니다.

Quickstart가 사용하는 전용 고객 키 999999의 상태를 Java에서 바꿔 보겠습니다. UpsertWriter는 Primary Key Table에서만 만들 수 있습니다. 실행할 때마다 서로 다른 금액을 만들고, 쓰기 전 값·첫 기록 뒤 조회 결과·같은 키로 바꿀 필드·변경 뒤 조회 결과를 차례로 출력합니다.

static BigDecimal randomBalanceExcept(BigDecimal previous) {
    BigDecimal amount;
    do {
        long cents = ThreadLocalRandom.current().nextLong(1_000, 100_000);
        amount = BigDecimal.valueOf(cents, 2);
    } while (amount.equals(previous));
    return amount;
}

static GenericRow profileRow(int customerId, String membership, BigDecimal balance) {
    return GenericRow.of(
        customerId,
        BinaryString.fromString("Quickstart User"),
        BinaryString.fromString(membership),
        Decimal.fromBigDecimal(balance, 15, 2)
    );
}

static void printProfile(String step, int customerId, String membership, BigDecimal balance) {
    System.out.printf("""
        %s
        {
          "customer_id": %d,
          "name": "Quickstart User",
          "membership": "%s",
          "account_balance": "%s"
        }
        %n""", step, customerId, membership, balance.toPlainString());
}

static void printProfile(String step, InternalRow row) {
    if (row == null) {
        throw new IllegalStateException("customer 999999 was not found");
    }
    printProfile(
        step,
        row.getInt(0),
        row.getString(2).toString(),
        row.getDecimal(3, 15, 2).toBigDecimal()
    );
}

int customerId = 999999;
BigDecimal beforeBalance = randomBalanceExcept(null);
BigDecimal afterBalance = randomBalanceExcept(beforeBalance);

try (Table profiles = connection.getTable(
        TablePath.of("demo", "customer_profile"))) {
    UpsertWriter writer = profiles.newUpsert().createWriter();
    Lookuper lookuper = profiles.newLookup().createLookuper();

    printProfile("[1/4] Initial value to write:", customerId, "java-client-before", beforeBalance);
    writer.upsert(profileRow(customerId, "java-client-before", beforeBalance)).get();
    writer.flush();

    printProfile("[2/4] Read the value after the first write:",
        lookuper.lookup(GenericRow.of(customerId)).get().getSingletonRow());

    System.out.printf("[3/4] Change fields for the same primary key:%n"
            + "membership=java-client-after%naccount_balance=%s%n", afterBalance);
    writer.upsert(profileRow(customerId, "java-client-after", afterBalance)).get();
    writer.flush();

    printProfile("[4/4] Read the value after the update:",
        lookuper.lookup(GenericRow.of(customerId)).get().getSingletonRow());
}

위 코드는 BigDecimal, ThreadLocalRandom, InternalRow을 사용합니다. upsert()는 비동기로 결과를 돌려주고, flush()는 앞서 보낸 쓰기가 서버의 확인 또는 오류에 도달할 때까지 기다립니다. 두 조회는 전체 행을 JSON 형태로 출력하므로 같은 customer_id에서 membershipaccount_balance가 바뀌고 행 수는 늘지 않는 것을 바로 확인할 수 있습니다.

키 조회가 반환하는 내부 행을 JSON으로 표시한다

Primary Key Table은 전체 이력을 훑는 대신 키를 기준으로 현재 상태를 찾습니다. 조회 키도 GenericRow로 만듭니다. 위 printProfile의 두 번째 오버로드처럼 결과가 없으면 getSingletonRow()null을 돌려주는 점을 별도로 처리해야 합니다. InternalRow의 문자열은 BinaryString이므로 toString()으로, DECIMAL(15, 2)getDecimal(3, 15, 2)으로 읽습니다.

LookupResult result = lookuper.lookup(GenericRow.of(customerId)).get();
printProfile("Current value:", result.getSingletonRow());

Lookuper도 스레드 안전하지 않습니다. 한 객체를 여러 요청 스레드가 공유하는 대신, 스레드별로 만들거나 호출을 직렬화해야 합니다.

Log Table에는 여러 이벤트를 추가하고 내용을 다시 확인한다

order_events에는 기본 키가 없습니다. 그러므로 같은 order_id를 다시 보내도 기존 행을 바꾸지 않으며, UpsertWriter가 아니라 AppendWriter를 사용합니다. DATE는 저수준 행 형식에서 epoch 이후 날짜 수인 int로 표현합니다. 네 이벤트를 writer에 모두 넣은 뒤 한 번만 flush()하고, 입력한 내용을 JSON으로 먼저 출력합니다.

record ExampleEvent(
        long orderId, int customerId, BigDecimal totalPrice,
        LocalDate orderedOn, String priority, String clerk) {}

static GenericRow eventRow(ExampleEvent event) {
    return GenericRow.of(
        event.orderId(), event.customerId(),
        Decimal.fromBigDecimal(event.totalPrice(), 15, 2),
        (int) event.orderedOn().toEpochDay(),
        BinaryString.fromString(event.priority()),
        BinaryString.fromString(event.clerk())
    );
}

static void printEvent(String step, ExampleEvent event) {
    System.out.printf("""
        %s
        {
          "order_id": %d,
          "customer_id": %d,
          "total_price": "%s",
          "ordered_on": "%s",
          "order_priority": "%s",
          "clerk": "%s"
        }
        %n""", step, event.orderId(), event.customerId(),
        event.totalPrice().toPlainString(), event.orderedOn(), event.priority(), event.clerk());
}

long orderIdBase = System.currentTimeMillis() * 10;
List<ExampleEvent> eventsToAppend = List.of(
    new ExampleEvent(orderIdBase, 999999, new BigDecimal("42.15"), LocalDate.of(2026, 8, 1), "low", "Java Client Clerk"),
    new ExampleEvent(orderIdBase + 1, 999999, new BigDecimal("85.60"), LocalDate.of(2026, 8, 2), "medium", "Java Client Clerk"),
    new ExampleEvent(orderIdBase + 2, 999999, new BigDecimal("129.90"), LocalDate.of(2026, 8, 3), "high", "Java Client Clerk"),
    new ExampleEvent(orderIdBase + 3, 999999, new BigDecimal("24.50"), LocalDate.of(2026, 8, 4), "low", "Java Client Clerk")
);
eventsToAppend.forEach(event -> printEvent("[1/3] Event queued for append:", event));

try (Table events = connection.getTable(TablePath.of("demo", "order_events"))) {
    AppendWriter writer = events.newAppend().createWriter();
    for (ExampleEvent event : eventsToAppend) {
        writer.append(eventRow(event));
    }
    writer.flush();
}

이 코드는 List를 사용합니다. Java Client 0.9.1의 AppendResult는 아직 저장 위치 offset을 제공하지 않습니다. 따라서 append 직후 “방금 쓴 네 건”을 offset으로 곧바로 다시 읽는 예제가 아닙니다. 소비자는 자신이 관리하는 시작 offset에서 버킷을 구독하고, 업무 식별자와 처리 위치를 별도로 관리해야 합니다.

Quickstart의 order_eventsbucket.num = 1인 비파티션 테이블이므로 버킷 0을 처음부터 구독할 수 있습니다. 처음부터 읽으면 이미 쌓인 주문도 함께 지나가므로, 이 예제에서는 방금 쓴 네 order_id를 찾아 실제 행 전체를 JSON으로 출력합니다. 무한 루프를 피하기 위해 15초 안에 모두 찾지 못하면 오류로 끝냅니다. 실제 소비자는 처음부터 반복해서 읽지 않고, 보관한 처리 위치에서 시작해야 합니다.

try (Table events = connection.getTable(TablePath.of("demo", "order_events"));
        LogScanner scanner = events.newScan().createLogScanner()) {
    scanner.subscribeFromBeginning(0);

    Set<Long> remaining = eventsToAppend.stream()
        .map(ExampleEvent::orderId)
        .collect(Collectors.toCollection(LinkedHashSet::new));
    long deadline = System.nanoTime() + Duration.ofSeconds(15).toNanos();
    while (!remaining.isEmpty() && System.nanoTime() < deadline) {
        ScanRecords records = scanner.poll(Duration.ofSeconds(1));
        for (ScanRecord record : records) {
            InternalRow row = record.getRow();
            if (remaining.remove(row.getLong(0))) {
                ExampleEvent read = new ExampleEvent(
                    row.getLong(0), row.getInt(1),
                    row.getDecimal(2, 15, 2).toBigDecimal(),
                    LocalDate.ofEpochDay(row.getInt(3)),
                    row.getString(4).toString(), row.getString(5).toString()
                );
                printEvent("[3/3] Read the appended event:", read);
            }
        }
    }
    if (!remaining.isEmpty()) {
        throw new IllegalStateException("did not read every appended event: " + remaining);
    }
}

위 코드는 Set, LinkedHashSet, Collectors를 사용합니다. 콘솔에는 네 입력 메시지와 스캔으로 다시 읽은 네 메시지가 모두 출력됩니다. 동일한 order_id는 Log Table에서 중복 방지 장치가 아니며, 이 예제의 시간 기반 ID는 실습에서만 사용합니다. 여러 버킷을 쓰는 실제 테이블에서는 각 버킷을 구독하고, 버킷별 시작 위치와 처리 완료 위치를 애플리케이션이 보관해야 합니다. 파티션 테이블은 subscribeFromBeginning(partitionId, bucket)처럼 파티션 ID도 함께 지정합니다.

GenericRow와 Typed API, 무엇을 고를까

Typed API는 POJO의 필드 이름과 타입을 테이블 스키마에 맞춰 변환합니다. order_events처럼 LocalDate, BigDecimal, String을 자연스럽게 쓰고 싶은 애플리케이션에는 이 쪽이 읽기 쉽습니다.

public class OrderEvent {
    public Long order_id;
    public Integer customer_id;
    public BigDecimal total_price;
    public LocalDate ordered_on;
    public String order_priority;
    public String clerk;

    public OrderEvent() {}
}

try (Table events = connection.getTable(TablePath.of("demo", "order_events"))) {
    TypedAppendWriter<OrderEvent> writer = events.newAppend()
        .createTypedWriter(OrderEvent.class);
    OrderEvent event = new OrderEvent();
    event.order_id = 1_000_000_002L;
    event.customer_id = 999999;
    event.total_price = new BigDecimal("42.00");
    event.ordered_on = LocalDate.of(2026, 8, 1);
    event.order_priority = "java-client";
    event.clerk = "Java Client Clerk";
    writer.append(event).get();
    writer.flush();
}
선택 기준GenericRow 기반 APIPOJO 기반 Typed API
값 표현Fluss 내부 형식과 필드 순서를 직접 다룸Java 필드 이름·타입을 스키마에 맞춤
적합한 경우변환 비용과 표현을 세밀하게 제어해야 할 때일반 애플리케이션 코드와 빠른 실습
주의점BinaryString, Decimal, 날짜의 내부 표현을 정확히 알아야 함POJO와 스키마가 어긋나면 변환에 실패하고 변환 비용이 추가됨

Typed API는 편의를 위해 내부 행과 POJO 사이를 변환합니다. 처리량과 지연 시간이 가장 중요한 경로에서는 GenericRow 같은 저수준 API를 검토할 수 있지만, 그 선택이 수명 주기 규칙을 바꾸지는 않습니다. Connection은 공유하고, Table·writer·scanner는 작업 범위에서 만듭니다.

실습을 마치기 전에 확인할 세 가지

  • Connection은 애플리케이션에서 공유하되 Table, Admin, Lookuper는 여러 스레드에서 함께 쓰지 않습니다.
  • customer_profileUpsertWriter와 키 조회를, order_eventsAppendWriter와 버킷 구독을 사용합니다.
  • Log Table의 소비 위치와 중복 처리 기준은 애플리케이션이 업무 식별자와 함께 관리합니다.

함께 읽기 좋은 글