Engineering Note
센서 이벤트에 빠진 현재 상태를 붙이는 방법: fluss-go 데이터 보강
fluss-go beta.10으로 센서 메타데이터 Primary Key Table과 측정값 Log Table을 만들고, 이벤트를 스캔하며 현재 상태를 조회해 결합합니다.
측정값 이벤트에는 온도와 습도만 있고, 센서 이름·설치 위치·정비 상태는 다른 테이블에 있다면 어떻게 처리할까요? 이벤트를 읽을 때 기본 키 테이블(Primary Key Table)에서 현재 센서 정보를 조회해 붙이면 됩니다. 별도 캐시나 외부 키-값 저장소를 두지 않아도 Fluss 안에서 최신 상태와 이벤트 이력을 함께 다룰 수 있습니다.
이 글은 Apache Fluss 공식 글 Apache Fluss Java Client: A Deep Dive의 IoT 데이터 보강 시나리오를 참고했습니다. 다만 Java 코드를 옮긴 글은 아닙니다. Apache Fluss 0.9.1-incubating과 fluss-go v0.1.0-beta.10을 기준으로, 데이터베이스·테이블·입력값·출력 형식을 새로 구성했습니다.
fluss-go는 Go 애플리케이션에서 Fluss의 데이터와 관리 API를 사용할 수 있게 하는 공개 베타 클라이언트입니다. 이 실습에서는 관리 API fadm으로 테이블을 만들고, 데이터 API fgo로 기본 키 테이블에 현재 상태를 기록한 뒤 로그 테이블(Log Table)의 이벤트를 읽고 보강합니다.
Konduo는 Fluss 같은 데이터 인프라를 관리·운영·모니터링하는 도구입니다. 이처럼 이벤트와 현재 상태를 함께 다루는 흐름을 이해해 두면, 운영 도구가 확인한 테이블·버킷 정보를 애플리케이션의 데이터 경로와 연결해 해석하기 쉬워집니다.
이 글에서 확인할 흐름
sensor_info 기본 키 테이블 (Primary Key Table)
sensor_id -> 이름, 위치, 현재 상태
sensor_readings 로그 테이블 (Log Table)
sensor_id, 측정 시각, 온도, 습도
LogScanner -> sensor_id 추출 -> Lookuper -> 보강된 측정값 출력
두 테이블의 역할을 섞지 않는 것이 핵심입니다. sensor_info는 같은 sensor_id를 다시 쓰면 현재 행이 바뀌는 기본 키 테이블입니다. sensor_readings는 새 측정값을 계속 쌓는 로그 테이블입니다. 읽는 쪽은 로그 테이블의 순서를 따르되, 각 이벤트의 sensor_id로 기본 키 테이블을 점 조회합니다.
데이터 보강이 필요한 이유
이벤트마다 센서 이름·설치 위치·정비 상태를 함께 넣으면, 상태가 바뀔 때 이미 기록한 이벤트와 새 이벤트의 정보가 쉽게 어긋납니다. 이 예제처럼 이벤트에는 센서 ID와 측정값을 남기고, 현재 상태는 기본 키 테이블에서 따로 관리하면 상태를 한곳에서 갱신할 수 있습니다.
로그 이벤트를 읽을 때 현재 상태를 조회해 결합하는 작업을 데이터 보강(enrichment)이라고 합니다. 대시보드·알림·실시간 분석처럼 현재의 센서 맥락이 필요한 소비자에게 특히 유용합니다. 이 방식은 과거 어느 시점의 상태를 정확히 되살리는 이력 조인이 아니라, 읽는 시점의 최신 상태를 붙이는 예제입니다.
이 흐름에서 Fluss를 쓰는 이유
일반적으로는 이벤트 로그를 메시징 시스템에 두고, 현재 상태는 별도의 키-값 저장소나 데이터베이스에 두어 애플리케이션이 두 시스템을 함께 읽습니다. 이 예제에서 Fluss는 로그 테이블의 순차 스캔과 기본 키 테이블의 점 조회를 같은 테이블 API와 클라이언트로 제공합니다. 따라서 이벤트·현재 상태의 데이터 경로와 테이블 메타데이터를 한 시스템 안에서 다룰 수 있습니다.
그 결과 애플리케이션은 보강을 위해 별도 저장소의 연결 설정·키 직렬화·조회 코드를 추가하지 않아도 됩니다. 다만 Fluss를 쓴다고 데이터 보강의 모든 문제가 사라지는 것은 아닙니다. 조회 실패 처리, 상태가 없는 센서의 처리, 그리고 과거 시점 상태가 필요한지 여부는 여전히 애플리케이션이 정해야 합니다.
이 실습에서는 다음 세 가지를 확인합니다.
- 센서 10개의 초기 상태를 기록한 뒤 두 센서 상태를 다시 기록해, 현재 상태 행이 10개로 유지되는지 확인합니다.
- 측정값 10건을 추가한 뒤, 방금 기록한 오프셋 범위만 다시 읽는지 확인합니다.
sensor_id2의 측정값에 나중에 갱신한OK상태가 붙는지 확인합니다.
실행 전 준비
Fluss 0.9.1-incubating 클러스터가 localhost:9123에서 실행 중이어야 합니다. Apache Fluss Flink Quickstart 환경을 그대로 사용해도 됩니다. 이 예제는 기존 demo 데이터베이스를 건드리지 않고 go_sensor_demo 데이터베이스와 두 테이블을 만듭니다.
전체 예제 파일을 내려받아 빈 디렉터리에서 실행합니다.
mkdir fluss-go-sensor-enrichment && cd fluss-go-sensor-enrichment
# 내려받은 sensor-enrichment.go 파일을 이 디렉터리에 둔다.
go mod init fluss-go-sensor-enrichment
go get github.com/pletorco/fluss-go@v0.1.0-beta.10
go mod tidy
go run sensor-enrichment.go
# Coordinator 주소가 다를 때
FLUSS_BOOTSTRAP=fluss.example:9123 go run sensor-enrichment.go
go mod tidy는 fgo와 fadm이 사용하는 의존성을 현재 소스의 import에 맞춰 기록합니다. fluss-go는 v1 이전의 공개 베타입니다. 이 예제는 beta.10의 GetTable, UpsertWriter, AppendWriter, Lookuper 이름을 사용하므로, beta.9 코드와 섞어 쓰지 않습니다. 또한 Fluss 0.9.1-incubating과 호환되도록 측정 시각 열의 TIMESTAMP(3) 정밀도를 명시합니다.
아래 코드는 실행 파일에서 핵심 흐름을 순서대로 발췌한 것입니다. 코드 조각만 이어 붙여 실행하지 말고, 실행할 때는 위에서 내려받은 전체 파일을 사용합니다. 본문은 각 API가 맡는 일을 설명하고, 전체 파일은 import·자료형·종료 처리까지 포함한 실행 가능한 형태를 제공합니다.
1. 관리 API로 두 테이블을 만든다
fadm.New는 이미 연결한 fgo.Client를 공유합니다. 따라서 관리 클라이언트를 닫는다고 별도 연결이 정리되는 구조가 아닙니다. 애플리케이션이 끝날 때 공유 fgo.Client를 닫습니다.
ctx := context.Background()
client, err := fgo.Open(
ctx,
fgo.WithBootstrapServers("localhost:9123"),
fgo.WithClientSoftware("sensor-enrichment", "0.1.0"),
)
if err != nil {
log.Fatal(err)
}
defer closeClient(client)
admin, err := fadm.New(client)
if err != nil {
log.Fatal(err)
}
if err := admin.CreateDatabase(ctx, "go_sensor_demo", fadm.DatabaseDescriptor{
Comment: "sensor-enrichment tutorial",
}, true); err != nil {
log.Fatal(err)
}
다음은 기본 키 테이블 정의입니다. sensor_id를 기본 키와 버킷 키로 지정합니다. BucketCount를 1로 둔 이유는 이 실습에서 방금 추가한 10개 이벤트의 오프셋 범위를 한 줄로 확인하기 위해서입니다. 처리량과 병렬성을 위한 운영 테이블이라면 예상 부하와 키 분포에 맞춰 버킷 수를 정해야 합니다.
sensorInfoPath := fgo.TablePath{Database: "go_sensor_demo", Table: "sensor_info"}
if err := admin.CreateTable(ctx, sensorInfoPath, fadm.TableDescriptor{
Comment: "current sensor metadata",
BucketCount: 1,
Schema: fgo.Schema{
Columns: []fgo.Column{
{Name: "sensor_id", Type: fgo.IntType},
{Name: "name", Type: fgo.StringType},
{Name: "location", Type: fgo.StringType},
{Name: "state", Type: fgo.StringType},
},
PrimaryKey: []string{"sensor_id"},
BucketKey: []string{"sensor_id"},
},
}, true); err != nil {
log.Fatal(err)
}
sensor_readings는 기본 키가 없는 로그 테이블입니다. 같은 sensor_id가 여러 번 와도 각 측정값은 별도 이벤트로 남습니다.
readingPath := fgo.TablePath{Database: "go_sensor_demo", Table: "sensor_readings"}
if err := admin.CreateTable(ctx, readingPath, fadm.TableDescriptor{
Comment: "append-only sensor readings",
BucketCount: 1,
Schema: fgo.Schema{
Columns: []fgo.Column{
{Name: "sensor_id", Type: fgo.IntType},
{
Name: "measured_at",
Type: fgo.TimestampType,
LogicalType: &fgo.LogicalType{
Root: "TIMESTAMP_WITHOUT_TIME_ZONE",
Precision: 3,
},
},
{Name: "temperature_c", Type: fgo.DoubleType},
{Name: "humidity_pct", Type: fgo.DoubleType},
},
BucketKey: []string{"sensor_id"},
},
}, true); err != nil {
log.Fatal(err)
}
ignoreIfExists를 true로 준 것은 실습을 다시 실행할 수 있게 하기 위해서입니다. 이미 같은 이름의 테이블이 있어도 스키마를 비교하거나 재생성하지는 않습니다. 운영 코드에서 스키마 변경을 처리할 때는 존재 여부만 무시하지 말고, 현재 정의를 읽어 검증하거나 별도의 마이그레이션 절차를 둬야 합니다.
2. 현재 센서 상태를 갱신한다
테이블을 만든 직후에는 GetTable로 서버의 테이블·스키마·버킷 정보를 읽습니다. UpsertWriter는 기본 키 테이블 전용 쓰기 객체입니다.
infos, err := client.GetTable(ctx, sensorInfoPath)
if err != nil {
log.Fatal(err)
}
inputInfos := []sensorInfo{
{1, "Roof temperature sensor", "roof", "OK"},
{2, "Lobby humidity sensor", "lobby", "ERROR"},
{3, "Server room sensor", "server-room", "MAINTENANCE"},
{4, "Warehouse pressure sensor", "warehouse", "OK"},
{5, "Conference room humidity sensor", "conference-room", "OK"},
{6, "Office 1 temperature sensor", "office-1", "LOW_BATTERY"},
{7, "Office 2 humidity sensor", "office-2", "OK"},
{8, "Lab pressure sensor", "lab", "ERROR"},
{9, "Parking pressure sensor", "parking", "OK"},
{10, "Backyard temperature sensor", "backyard", "OK"},
{2, "Lobby humidity sensor", "lobby", "OK"}, // 상태 변경
{8, "Lab pressure sensor", "lab", "CALIBRATING"}, // 상태 변경
}
writer, err := client.NewUpsertWriter(
ctx,
infos,
fgo.WithUpsertBatchLimits(1<<20, len(inputInfos)),
)
if err != nil {
log.Fatal(err)
}
defer closeUpsertWriter(ctx, writer)
for _, info := range inputInfos {
result := writer.Upsert(ctx, fgo.Row{
info.SensorID, info.Name, info.Location, info.State,
}).Await(ctx)
if result.Err != nil {
log.Fatal(result.Err)
}
}
if err := writer.Flush(ctx); err != nil {
log.Fatal(err)
}
Await는 각 쓰기의 결과를 확인하고, Flush는 호출 전에 writer가 받은 쓰기가 성공 또는 오류라는 최종 상태에 이를 때까지 기다립니다. 이 예제는 센서 10개의 초기 상태를 쓴 뒤 sensor_id 2와 8을 다시 씁니다. 따라서 입력은 12건이지만 기본 키 테이블의 최종 행 수는 10건이며, 두 센서에는 각각 OK, CALIBRATING 상태가 남습니다. 모든 열을 전달하므로 같은 키의 현재 행이 새 값으로 갱신됩니다.
3. 측정값을 추가하고 읽을 범위를 정한다
이벤트 테이블에는 AppendWriter를 사용합니다. 10건을 모두 writer에 넣은 뒤 한 번 Flush하고, 응답의 버킷과 오프셋을 보관합니다.
inputReadings := []sensorReading{
{1, time.Date(2026, time.August, 4, 9, 0, 0, 0, time.UTC), 24.1, 43.2},
{2, time.Date(2026, time.August, 4, 9, 15, 0, 0, time.UTC), 22.7, 52.8},
{3, time.Date(2026, time.August, 4, 9, 30, 0, 0, time.UTC), 20.4, 38.5},
{4, time.Date(2026, time.August, 4, 9, 45, 0, 0, time.UTC), 18.9, 48.1},
{5, time.Date(2026, time.August, 4, 10, 0, 0, 0, time.UTC), 23.5, 46.3},
{6, time.Date(2026, time.August, 4, 10, 15, 0, 0, time.UTC), 21.8, 44.9},
{7, time.Date(2026, time.August, 4, 10, 30, 0, 0, time.UTC), 22.1, 47.5},
{8, time.Date(2026, time.August, 4, 10, 45, 0, 0, time.UTC), 20.7, 49.2},
{9, time.Date(2026, time.August, 4, 11, 0, 0, 0, time.UTC), 19.6, 51.7},
{10, time.Date(2026, time.August, 4, 11, 15, 0, 0, time.UTC), 25.0, 41.8},
}
readings, err := client.GetTable(ctx, readingPath)
if err != nil {
log.Fatal(err)
}
writer, err := client.NewAppendWriter(
ctx,
readings,
fgo.WithAppendBatchLimits(1<<20, len(inputReadings)),
fgo.WithAppendBatchTimeout(5*time.Millisecond),
)
if err != nil {
log.Fatal(err)
}
defer closeAppendWriter(ctx, writer)
futures := make([]*fgo.WriteFuture, len(inputReadings))
for i, reading := range inputReadings {
futures[i] = writer.Append(ctx, fgo.Row{
reading.SensorID, reading.MeasuredAt,
reading.TemperatureC, reading.HumidityPct,
})
}
if err := writer.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 || !results[i].OffsetKnown {
log.Fatalf("append %d: %v", i, results[i].Err)
}
}
first, last := results[0], results[len(results)-1]
Append가 성공했더라도 오프셋이 확인되지 않았으면 그 위치를 읽기 시작점으로 삼을 수 없습니다. 이 코드는 모든 WriteFuture의 결과를 확인한 뒤 첫·마지막 오프셋을 기억합니다. 전체 예제는 여기에 더해 단일 버킷에서 오프셋 10개가 연속인지 검증합니다. 이 검증은 바로 다음 단계에서 다른 실행이 넣은 이벤트를 섞어 읽지 않기 위한 것입니다.
4. 이벤트에 현재 센서 상태를 붙인다
LogScanner는 이벤트의 오프셋을 따라 읽고, Lookuper는 sensor_id로 현재 센서 상태를 찾습니다. 아래 코드는 스캐너가 돌려준 배치를 처리한 뒤 반드시 Release합니다.
lookuper, err := client.NewLookuper(ctx, infos, fgo.WithLookupBatchLimits(100, 4))
if err != nil {
log.Fatal(err)
}
defer closeLookuper(lookuper)
scanner, err := client.NewLogScanner(
ctx,
readings,
fgo.AtOffset(first.BaseOffset),
fgo.WithScanRowLimit(int64(len(inputReadings))),
fgo.WithScanStoppingOffsets(map[int32]int64{
first.Bucket: last.BaseOffset + 1,
}),
)
if err != nil {
log.Fatal(err)
}
defer closeScanner(scanner)
enriched := make([]enrichedReading, 0, len(inputReadings))
for !scanner.Done() {
batch, err := scanner.Poll(ctx)
if err != nil {
log.Fatal(err)
}
for _, record := range batch.Records {
reading := readingFromRow(record.Record.Value)
lookups := lookuper.Lookup(ctx, fgo.PrimaryKey{reading.SensorID})
if len(lookups) != 1 {
log.Fatalf("sensor %d: expected one lookup result, got %d", reading.SensorID, len(lookups))
}
lookup := lookups[0]
switch {
case errors.Is(lookup.Err, fgo.ErrNotFound):
log.Fatalf("sensor %d has no current metadata", reading.SensorID)
case lookup.Err != nil:
log.Fatal(lookup.Err)
}
info := infoFromRow(lookup.Row)
enriched = append(enriched, enrichedReading{
SensorID: reading.SensorID, MeasuredAt: reading.MeasuredAt,
TemperatureC: reading.TemperatureC, HumidityPct: reading.HumidityPct,
SensorName: info.Name, Location: info.Location, State: info.State,
})
}
batch.Release()
}
printJSON("[4/4] Scan readings and enrich each with current sensor metadata:", enriched)
전체 파일은 각 단계의 입력과 결과를 JSON으로 출력합니다. 마지막에는 다음과 같이 이벤트와 현재 상태가 결합됩니다.
{
"sensor_id": 2,
"measured_at": "2026-08-04T09:15:00Z",
"temperature_c": 22.7,
"humidity_pct": 52.8,
"sensor_name": "Lobby humidity sensor",
"location": "lobby",
"state": "OK"
}
sensor_id 2는 처음에 ERROR 상태로 기록했지만, 이벤트를 추가하기 전에 OK로 다시 기록했습니다. 따라서 이 출력은 이벤트가 기록된 시점의 상태가 아니라, 조회한 시점의 현재 상태를 보여 줍니다.
이 예제가 보장하지 않는 것
이 흐름은 이벤트를 읽는 시점의 현재 상태를 붙입니다. 측정값이 기록된 당시의 센서 상태를 재현하는 시간 기준 결합은 아닙니다. 예를 들어 이벤트를 쓴 뒤 센서 위치가 바뀌고 나서 스캔하면, 조회 결과에는 바뀐 위치가 들어갈 수 있습니다.
과거 시점의 정확한 상태가 필요하다면 이벤트에 필요한 메타데이터를 함께 넣거나, 버전·유효 시점이 있는 상태 이력을 별도로 설계해야 합니다. 또한 이 예제는 읽은 행마다 한 번씩 조회해 흐름을 드러냅니다. 높은 처리량의 애플리케이션에서는 여러 키를 모아 Lookuper에 전달하고, 누락된 키·부분 실패·타임아웃을 각각 다루는 경계를 정해야 합니다.
정리
Fluss의 로그 테이블은 시간순 사실을 남기고, 기본 키 테이블은 현재 상태를 제공합니다. AppendWriter, LogScanner, Lookuper를 각각의 역할에 맞게 연결하면 Go 애플리케이션도 이 두 모델을 직접 결합할 수 있습니다. 다만 “현재 상태를 붙이는 일”과 “과거 시점의 상태를 복원하는 일”은 다르므로, 필요한 시간 의미를 먼저 정해야 합니다.
함께 읽기 좋은 글
- Go로 Apache Fluss 테이블 읽고 쓰기: fluss-go 공개 베타 - beta.10 API 이름과 Primary Key Table·Log Table의 기본 읽기·쓰기 흐름을 먼저 확인합니다.
- Connection은 하나, Table은 스레드마다: Apache Fluss Java Client 실습 - 같은 두 테이블 모델을 Java Client의 객체 수명 주기 관점에서 비교합니다.
- Apache Fluss 첫 실행: Flink Quickstart에서 확인할 것들 - 로컬 Fluss 환경을 먼저 준비하고 테이블 모델을 관찰합니다.