Engineering Note

Go로 Apache Fluss 테이블 읽고 쓰기: fluss-go 공개 베타

Apache Fluss Flink Quickstart에서 만든 demo.customer_profile과 demo.order_events를 fluss-go 공개 베타로 갱신·조회·추가·스캔하는 방법을 정리합니다.

2026년 8월 1일 · Pletor Engineering apache-flussgoclient-libraryprotocolopen-source

Apache Fluss Flink Quickstart에서는 Flink SQL로 demo.customer_profiledemo.order_events를 만들고, 같은 키의 최신 상태와 계속 쌓이는 이벤트 이력이 어떻게 다른지 확인했습니다. 이번에는 그 같은 테이블을 Go 프로그램에서 직접 다룹니다.

fluss-go는 Apache Fluss용 Go 클라이언트 라이브러리입니다. 현재 v0.1.0-beta.10Apache Fluss 0.9.1-incubating 지원하는 공개 베타입니다. 데이터 작업 API fgo, 관리 API fadm, 바이너리 프로토콜 메시지 fmsg가 이 공개 베타를 이룹니다. 이 글에서는 데이터 API를 통해 고객 상태와 주문 이벤트를 다뤄 보면서, 라이브러리 전체의 계층과 검증 범위를 함께 살펴봅니다.

Konduo는 Fluss 클러스터를 관리·운영·모니터링하는 도구입니다. fluss-go는 Go 기반 Konduo Fluss 플러그인이 테이블, 버킷, 클러스터 상태를 읽고 필요한 운영 작업을 수행할 때 쓰는 기반 라이브러리입니다. 이 글의 예제는 그 라이브러리를 애플리케이션 코드에서 먼저 확인하는 방법이기도 합니다.

청록색 최신 상태 격자와 파란색 이벤트 레일을 중앙의 Go 클라이언트가 연결하는 등각형 일러스트
하나의 Go 클라이언트로 최신 상태와 계속 쌓이는 이벤트 이력을 함께 다룹니다.

시작하기 전에: Quickstart의 테이블을 그대로 쓴다

먼저 Flink Quickstart의 테이블 생성과 데이터 적재 단계를 마칩니다. 이 글은 새 데이터베이스나 별도 스키마를 만들지 않습니다. Go 코드가 여는 대상은 다음 두 테이블입니다.

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

customer_profilecustomer_idINT이므로 Go 코드에서는 int32를 사용합니다. order_events.order_idBIGINT라서 int64입니다. DECIMAL(15, 2) 값은 *big.Rat, DATE 값은 time.Time으로 전달합니다. 스키마에 맞지 않는 Go 값을 넣으면 writer가 요청을 만들기 전에 오류를 반환합니다.

코드를 실행하기 전에 go version으로 Go 버전을 확인합니다. Go 1.25 계열은 1.25.12 이상, Go 1.26 계열은 1.26.5 이상이 필요합니다. Quickstart에서 만든 두 테이블의 스키마도 위 정의와 같은지 확인합니다.

예제 전체에는 다음 import가 필요합니다.

import (
    "context"
    "encoding/json"
    "errors"
    "fmt"
    "log"
    "math/big"
    "math/rand/v2"
    "time"

    "github.com/pletorco/fluss-go/pkg/fgo"
)

아래 예제는 파일로도 받을 수 있습니다. Quickstart를 실행한 뒤, 빈 디렉터리에서 Go 모듈을 만들고 의존성을 받습니다. 기본 Coordinator 주소는 localhost:9123이며, 다른 주소를 쓴다면 FLUSS_BOOTSTRAP 환경 변수로 지정합니다.

mkdir fluss-go-example && cd fluss-go-example
go mod init fluss-go-example
go get github.com/pletorco/fluss-go/pkg/fgo@v0.1.0-beta.10

# Go 1.25.12 이상 또는 Go 1.26.5 이상이 필요하다.
# 내려받은 파일을 이 디렉터리에 둔 뒤 실행한다.
go run primary-key-table.go
go run log-table.go

# Coordinator 주소가 다를 때
FLUSS_BOOTSTRAP=fluss.example:9123 go run primary-key-table.go

공개 베타의 클라이언트 계층

공개 베타의 중심에는 Coordinator·TabletServer 연결, 버전 협상, 인증, 메타데이터와 버킷 라우팅을 관리하는 fgo.Client가 있습니다. 데이터 API fgo와 관리 API fadm은 이 연결 기반을 공유합니다. beta.10은 실험적 공개 API의 이름을 Apache Fluss 용어에 맞췄으며, 이전 이름의 호환 별칭은 제공하지 않습니다. 아래 코드는 beta.10 이름으로 두 테이블을 읽습니다. beta.9에서 추가된 Coordinator 연결 교체 뒤 논리적인 fgo.Client를 유지하는 동작도 그대로 이어집니다.

Go 애플리케이션이 fgo 데이터 API와 fadm 관리 API를 사용하고, 두 API가 fmsg 프로토콜 메시지와 내부 전송 계층을 거쳐 Apache Fluss Coordinator와 TabletServer에 연결되는 구조도
`fgo`와 `fadm`은 같은 클라이언트 연결을 공유합니다. 이 글의 실습은 `fgo`를 사용하지만, 공개 베타의 범위는 관리 API와 프로토콜·전송 계층까지 포함합니다.
ctx := context.Background()

client, err := fgo.Open(
    ctx,
    fgo.WithBootstrapServers("coordinator.example:9123"),
    fgo.WithClientSoftware("quickstart-go", "0.1.0"),
)
if err != nil {
    log.Fatal(err)
}
defer client.Close()

profiles, err := client.GetTable(ctx, fgo.TablePath{
    Database: "demo",
    Table:    "customer_profile",
})
if err != nil {
    log.Fatal(err)
}

events, err := client.GetTable(ctx, fgo.TablePath{
    Database: "demo",
    Table:    "order_events",
})
if err != nil {
    log.Fatal(err)
}

GetTable은 이름만 담은 핸들을 만들지 않습니다. 서버의 테이블·스키마 정보를 읽고, 요청을 어느 버킷과 TabletServer로 보낼지 결정하는 정보를 채웁니다. Quickstart에서 스키마를 바꿨다면 새 스키마가 보일 때까지 테이블을 다시 읽고, 그 새 핸들로 읽기·쓰기 객체를 만들어야 합니다.

Primary Key Table: Quickstart 고객 상태를 Go에서 갱신하고 조회한다

Quickstart는 실습 전용 고객 키 999999를 사용해 같은 키의 두 번째 INSERT가 최신 상태를 바꾼다는 점을 보여 줍니다. 이 예제도 같은 키에 먼저 go-client-before를 기록하고 조회한 뒤, go-client-after로 다시 upsert해 조회합니다. 그래서 이전 실행 상태와 관계없이 콘솔에서 현재 상태가 바뀌는 과정을 확인할 수 있습니다. 이 Quickstart 테이블에는 table.merge-engine이 없으므로 MergeModeOverwrite를 지정하지 않고 기본 병합 모드를 사용합니다.

lookup, err := client.NewLookuper(
    ctx,
    profiles,
    fgo.WithLookupBatchLimits(100, 4),
)
if err != nil {
    log.Fatal(err)
}
defer lookup.Close()

writer, err := client.NewUpsertWriter(ctx, profiles, fgo.WithUpsertBatchLimits(1<<20, 500))
if err != nil {
    log.Fatal(err)
}
defer writer.Close(ctx)

type customerProfile struct {
    CustomerID  int32  `json:"customer_id"`
    Name        string `json:"name"`
    Membership  string `json:"membership"`
    AccountBalance string `json:"account_balance"` // decimal은 정밀도를 보존하도록 문자열로 출력한다.
}

writeProfile := func(profile customerProfile) {
    accountBalance, ok := new(big.Rat).SetString(profile.AccountBalance)
    if !ok {
        log.Fatalf("invalid account_balance: %q", profile.AccountBalance)
    }
    result := writer.Upsert(ctx, fgo.Row{
        profile.CustomerID, profile.Name, profile.Membership, accountBalance,
    }).Await(ctx)
    if result.Err != nil {
        log.Fatal(result.Err)
    }
    if err := writer.Flush(ctx); err != nil {
        log.Fatal(err)
    }
}
showProfile := func(step string) {
    result := lookup.Lookup(ctx, fgo.PrimaryKey{int32(999999)})[0]
    if errors.Is(result.Err, fgo.ErrNotFound) {
        log.Printf("%s customer was not found", step)
    } else if result.Err != nil {
        log.Fatal(result.Err)
    } else {
        printJSON(step, profileFromRow(result.Row))
    }
}

before := customerProfile{int32(999999), "Quickstart User", "go-client-before", randomAmount("")}
after := before
after.Membership = "go-client-after"
after.AccountBalance = randomAmount(before.AccountBalance)

printJSON("[1/4] Initial value to write:", before)
writeProfile(before)
showProfile("[2/4] Read the value after the first write:")
log.Printf("[3/4] Change fields for the same primary key:\nmembership=%s\naccount_balance=%s", after.Membership, after.AccountBalance)
writeProfile(after)
showProfile("[4/4] Read the value after the update:")
func randomAmount(except string) string {
    for {
        cents := 1_000 + rand.IntN(99_000) // 10.00 ~ 999.99
        amount := fmt.Sprintf("%d.%02d", cents/100, cents%100)
        if amount != except {
            return amount
        }
    }
}
func profileFromRow(row fgo.Row) customerProfile {
    accountBalance, ok := row[3].(*big.Rat)
    if !ok {
        log.Fatalf("unexpected account_balance type: %T", row[3])
    }
    return customerProfile{
        CustomerID: row[0].(int32), Name: row[1].(string),
        Membership: row[2].(string), AccountBalance: accountBalance.FloatString(2),
    }
}

func printJSON(step string, value customerProfile) {
    encoded, err := json.MarshalIndent(value, "", "  ")
    if err != nil {
        log.Fatal(err)
    }
    log.Printf("%s\n%s", step, encoded)
}

1단계는 쓰기 전에 입력할 값을 먼저 출력합니다. randomAmount는 매 실행마다 10.00~999.99 범위의 서로 다른 금액을 두 개 만들어, 두 번째 upsert에서 membershipaccount_balance가 함께 바뀌도록 합니다. profileFromRowprintJSON은 각각 조회 행을 위 JSON 구조로 바꾸고 보기 좋게 출력하는 보조 함수입니다. account_balance는 실습의 소수 정밀도를 그대로 보여 주기 위해 문자열로 출력합니다. Lookuper는 키와 버킷을 이용해 현재 상태를 찾으므로 모든 행을 훑은 뒤 필터링하지 않습니다. 실행하면 다음 순서가 출력됩니다.

[1/4] Initial value to write:
{
  "customer_id": 999999,
  "name": "Quickstart User",
  "membership": "go-client-before",
  "account_balance": "418.27"
}
[2/4] Read the value after the first write:
{
  "customer_id": 999999,
  "name": "Quickstart User",
  "membership": "go-client-before",
  "account_balance": "418.27"
}
[3/4] Change fields for the same primary key:
membership=go-client-after
account_balance=762.94
[4/4] Read the value after the update:
{
  "customer_id": 999999,
  "name": "Quickstart User",
  "membership": "go-client-after",
  "account_balance": "762.94"
}

두 조회에서 키는 같고 membershipaccount_balance가 바뀝니다. 같은 customer_idUpsert를 여러 번 실행해도 행 수가 늘지 않고 현재 값만 바뀝니다.

Lookup은 여러 키를 한 번에 받을 수 있으며, 내부적으로 버킷별 요청으로 묶습니다. 반환값은 입력 키와의 대응을 유지하므로, 찾지 못한 키와 서버·전송 오류를 구분할 수 있습니다. ErrNotFound는 정상적인 데이터 결과일 수 있지만, 타임아웃은 그렇지 않습니다.

Log Table: 주문 이벤트 네 건을 함께 추가하고 다시 읽는다

order_events에는 기본 키가 없습니다. 그러므로 고객 999999의 이벤트를 새로 기록할 때는 UpsertWriter가 아니라 AppendWriter를 사용합니다. 아래 예제는 Quickstart에서 넣은 행을 바꾸지 않고, 새 주문 이벤트 네 건을 한 번에 큐에 넣은 뒤 그 내용을 JSON으로 출력합니다.

appendWriter, err := client.NewAppendWriter(
    ctx,
    events,
    fgo.WithAppendBatchLimits(1<<20, 500),
    fgo.WithAppendBatchTimeout(5*time.Millisecond),
)
if err != nil {
    log.Fatal(err)
}
defer appendWriter.Close(ctx)

orderIDBase := time.Now().UnixNano()
eventsToAppend := []orderEvent{
    {orderIDBase, 999999, "42.15", "2026-08-01", "low", "Go Client Clerk"},
    {orderIDBase + 1, 999999, "85.60", "2026-08-02", "medium", "Go Client Clerk"},
    {orderIDBase + 2, 999999, "129.90", "2026-08-03", "high", "Go Client Clerk"},
    {orderIDBase + 3, 999999, "24.50", "2026-08-04", "low", "Go Client Clerk"},
}
printJSON("[1/3] Queue these four new events:", eventsToAppend)

futures := make([]*fgo.WriteFuture, len(eventsToAppend))
for i, event := range eventsToAppend {
    futures[i] = appendWriter.Append(ctx, event.row())
}
if err := appendWriter.Flush(ctx); err != nil {
    log.Fatal(err)
}

results := make([]fgo.WriteResult, len(futures))
for i, future := range futures {
    results[i] = future.Await(ctx)
    if results[i].Err != nil {
        log.Fatal(results[i].Err)
    }
    if !results[i].OffsetKnown {
        log.Fatal("append succeeded but did not return an offset")
    }
}
first := results[0]
for i, result := range results {
    expectedOffset := first.BaseOffset + int64(i)
    if result.Bucket != first.Bucket || result.BaseOffset != expectedOffset {
        log.Fatalf("events do not occupy one contiguous range")
    }
}
log.Printf("[2/3] Fluss stored four events in bucket=%d at offsets %d through %d", first.Bucket, first.BaseOffset, results[len(results)-1].BaseOffset)

Append는 완료를 기다리기 전에 모두 writer에 넣고, Flush 한 번으로 완료를 기다립니다. Quickstart는 버킷이 하나이므로 이 네 이벤트의 순서가 유지됩니다. 예제는 네 결과의 버킷과 offset이 연속인지 확인한 뒤, 그 구간의 끝 offset을 스캐너의 종료 지점으로 지정합니다. 그래서 다른 이벤트가 섞여도 네 건을 성공으로 오인하지 않습니다. Append가 성공하고 offset이 확인된 경우에만 BaseOffset으로 기록 위치를 확인할 수 있습니다. 같은 order_idcustomer_id를 다시 보내더라도 Log Table은 기존 행을 대체하지 않습니다. 이 코드를 반복 실행하면 이벤트가 매번 네 행씩 늘어납니다. 애플리케이션에서 중복을 막아야 한다면 안정적인 이벤트 ID와 소비자 측 멱등성 규칙을 별도로 둬야 합니다.

방금 추가한 네 이벤트를 offset에서 읽고 내용을 확인한다

Quickstart의 order_eventsbucket.num = 1인 테이블입니다. 따라서 첫 Append가 돌려준 BaseOffset에서 네 건을 읽으면, 방금 넣은 이벤트와 그 내용을 바로 확인할 수 있습니다.

scanner, err := client.NewLogScanner(
    ctx,
    events,
    fgo.AtOffset(first.BaseOffset),
    fgo.WithScanRowLimit(int64(len(eventsToAppend))),
    fgo.WithScanStoppingOffsets(map[int32]int64{
        first.Bucket: results[len(results)-1].BaseOffset + 1,
    }),
)
if err != nil {
    log.Fatal(err)
}
defer scanner.Close()

readEvents := make([]orderEvent, 0, len(eventsToAppend))
for !scanner.Done() {
    batch, err := scanner.Poll(ctx)
    if err != nil {
        log.Fatal(err)
    }
    for _, record := range batch.Records {
        readEvents = append(readEvents, orderEventFromRow(record.Record.Value))
    }
    batch.Release()
}
if len(readEvents) != len(eventsToAppend) {
    log.Fatalf("read %d events, want %d", len(readEvents), len(eventsToAppend))
}
printJSON("[3/3] Read the appended events:", readEvents)

orderEvent, event.row, orderEventFromRow, printJSON은 내려받는 전체 예제에 들어 있습니다. 소수점 정밀도를 유지하기 위해 total_price는 JSON에서 문자열로 표시합니다. 실행하면 입력한 네 건과 읽은 네 건이 같은 형태로 출력되어, 실제 내용까지 확인할 수 있습니다.

[1/3] Queue these four new events:
[
  {"order_id": 1760000000000000000, "customer_id": 999999, "total_price": "42.15", "ordered_on": "2026-08-01", "order_priority": "low", "clerk": "Go Client Clerk"},
  {"order_id": 1760000000000000001, "customer_id": 999999, "total_price": "85.60", "ordered_on": "2026-08-02", "order_priority": "medium", "clerk": "Go Client Clerk"},
  {"order_id": 1760000000000000002, "customer_id": 999999, "total_price": "129.90", "ordered_on": "2026-08-03", "order_priority": "high", "clerk": "Go Client Clerk"},
  {"order_id": 1760000000000000003, "customer_id": 999999, "total_price": "24.50", "ordered_on": "2026-08-04", "order_priority": "low", "clerk": "Go Client Clerk"}
]
[2/3] Fluss stored four events in bucket=0 at offsets 100 through 103
[3/3] Read the appended events:
[
  {"order_id": 1760000000000000000, "customer_id": 999999, "total_price": "42.15", "ordered_on": "2026-08-01", "order_priority": "low", "clerk": "Go Client Clerk"},
  {"order_id": 1760000000000000001, "customer_id": 999999, "total_price": "85.60", "ordered_on": "2026-08-02", "order_priority": "medium", "clerk": "Go Client Clerk"},
  {"order_id": 1760000000000000002, "customer_id": 999999, "total_price": "129.90", "ordered_on": "2026-08-03", "order_priority": "high", "clerk": "Go Client Clerk"},
  {"order_id": 1760000000000000003, "customer_id": 999999, "total_price": "24.50", "ordered_on": "2026-08-04", "order_priority": "low", "clerk": "Go Client Clerk"}
]

AtOffset은 해당 offset을 포함해 읽는 시작점을 뜻합니다. 실제 애플리케이션은 Earliest(), Latest(), 특정 타임스탬프, 명시한 offset을 선택할 수 있습니다. 여러 버킷을 쓰는 테이블에서는 버킷별 시작·종료 offset을 별도로 관리해야 합니다. Poll은 새 레코드를 기다릴 수 있으므로, 모든 호출에는 취소 가능한 context.Context가 필요합니다. beta.9에서는 서버의 바이트 제한 때문에 응답 마지막의 레코드 배치가 불완전해져도, 앞서 완전하게 받은 행·Arrow 배치를 보존하고 다음 fetch를 마지막 완전 offset부터 이어 갑니다. 유효한 응답을 잘못된 배치로 처리하지 않기 위한 변경입니다.

같은 Quickstart 테이블에서 확인할 차이

Flink SQL과 Go API는 같은 Fluss 테이블을 바라봅니다. 다만 테이블 모델이 다르므로 확인해야 하는 결과도 다릅니다.

작업customer_profile Primary Key Tableorder_events Log Table
쓰기customer_id를 기준으로 upsert 또는 delete새 주문 이벤트를 append
읽기키 또는 키 접두사 조회, 현재 상태 스캔버킷의 offset을 따라 순차 스캔
성공 확인키의 최신 행 또는 변경 결과기록된 버킷·offset
반복 실행같은 키의 최신 값 갱신이벤트 한 행 추가

이 차이를 하나의 Put·Get API로 뭉개지 않는 것이 fluss-go의 선택입니다. Fluss의 테이블 모델을 따라 API도 달라져야 잘못된 읽기와 재시도 정책을 줄일 수 있습니다.

시간 초과 뒤에는 자동 재전송하지 않는다

쓰기 성공 응답을 받지 못했다고 해서 서버가 쓰기를 거절한 것은 아닙니다. 네트워크 오류나 호출 취소 뒤에는 서버가 요청을 반영하고 응답만 잃었을 수 있습니다.

fluss-go는 안전한 읽기 요청에는 제한된 재시도를 적용할 수 있지만, 결과가 모호한 쓰기를 자동으로 다시 보내지 않습니다. 해당 버킷의 쓰기 객체는 ErrWriterState를 반환합니다. 호출자는 이전 쓰기 묶음의 결과를 확인하고, 애플리케이션의 멱등성 보장 방식에 따라 새 쓰기 객체를 만들거나 별도의 확인 절차를 거쳐야 합니다.

Flush(ctx)는 호출 전까지 쓰기 객체가 받아들인 쓰기가 성공 또는 오류라는 최종 결과에 도달할 때까지 기다립니다. Close(ctx)도 남은 쓰기를 flush한 뒤 리소스를 정리합니다. 이 수명 주기를 명확히 두지 않으면 프로그램 종료 직전에 버퍼에 남은 이벤트를 잃거나, 모호한 결과를 성공처럼 기록할 수 있습니다. fgo.Client.Close()는 beta.9에서 명시적으로 종결 상태가 되며 여러 번 호출해도 안전합니다. 따라서 Client는 모든 테이블·리더·라이터 작업을 끝낸 뒤에만 닫아야 합니다.

공개 베타의 구현·검증 범위

코드 예제는 기능이 존재한다는 출발점일 뿐입니다. fluss-go 공개 베타는 다음 방식으로 0.9.1 호환성을 확인합니다.

  • pkg/fmsg의 메시지는 Fluss 0.9.1의 FlussApi.proto, ApiKeys.java, Errors.java에서 생성합니다. 원본 파일의 SHA-256을 기록해 입력 변경을 검토합니다.
  • Java 호환 바이트 픽스처로 프레임, 행·KV·Log·Arrow 배치, 버킷 해시를 확인합니다.
  • 고정 다이제스트의 공식 Fluss 0.9.1 이미지에서 일반 연결, SASL PLAIN, 테이블 관리, append·scan, KV·lookup, Coordinator 연결 교체를 통합 검증합니다.
  • 행·Arrow 배치가 버킷별 fetch 바이트 제한을 넘어 끝부분이 잘린 응답도 단위·통합 테스트로 확인합니다.
  • 제한된 시간과 작업 수의 신뢰성 실행에서 취소, 잘린 연결, TabletServer 재시작 뒤에도 확인된 결과와 리소스 경계를 검사합니다.

이 검증은 Fluss 서버 성능 벤치마크가 아니며, 이후 Fluss 버전이나 모든 클라우드 저장소 조합을 보장하지 않습니다. 현재 v0.1.0-beta.10은 v1 이전의 공개 베타입니다. beta.9에서 올릴 때는 공개 API 이름이 바뀌므로 공식 마이그레이션 표를 먼저 확인합니다. 운영 코드에서는 go get ...@latest 대신 검토한 태그를 고정하고, 서버 버전도 0.9.1-incubating인지 확인해야 합니다. 선택 어댑터 모듈을 함께 쓴다면 루트 모듈과 어댑터 모두 beta.10 태그로 맞춥니다.

go get github.com/pletorco/fluss-go/pkg/fgo@v0.1.0-beta.10

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

  • Primary Key Table에는 UpsertWriterLookuper, Log Table에는 AppendWriterLogScanner를 사용합니다.
  • 쓰기 응답이 모호하게 끝나면 자동으로 다시 보내지 말고, event_id 같은 업무 식별자와 조회 절차로 먼저 결과를 확인합니다.

함께 읽기 좋은 글