From 77a3b3e86b06504bdf37e35d83a192806a95a63e Mon Sep 17 00:00:00 2001 From: SangTran-127 Date: Sat, 8 Aug 2026 11:31:29 +0700 Subject: initialize goflink module with core stream record implementation --- .gitignore | 1 + core/record.go | 40 ++++++++++++++++++++++++++++++++++++++++ go.mod | 3 +++ main.go | 13 +++++++++++++ 4 files changed, 57 insertions(+) create mode 100644 .gitignore create mode 100644 core/record.go create mode 100644 go.mod create mode 100644 main.go diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..62c8935 --- /dev/null +++ b/.gitignore @@ -0,0 +1 @@ +.idea/ \ No newline at end of file diff --git a/core/record.go b/core/record.go new file mode 100644 index 0000000..cb160b2 --- /dev/null +++ b/core/record.go @@ -0,0 +1,40 @@ +package core + +type RecordType byte + +const ( + RecordData RecordType = iota + RecordWatermark + RecordBarrier +) + +type StreamRecord[T any] struct { + Value T + Key string // for keyBy stateful processing + Timestamp int64 // epoch ms time + BarrierID uint64 // ID checkpoint barrier + Type RecordType +} + +func NewDataRecord[T any](key string, val T, timestamp int64) *StreamRecord[T] { + return &StreamRecord[T]{ + Value: val, + Key: key, + Timestamp: timestamp, + Type: RecordData, + } +} + +func NewWatermarkRecord[T any](timestamp int64) *StreamRecord[T] { + return &StreamRecord[T]{ + Type: RecordWatermark, + Timestamp: timestamp, + } +} + +func NewBarrierRecord[T any](barrierID uint64) *StreamRecord[T] { + return &StreamRecord[T]{ + Type: RecordBarrier, + BarrierID: barrierID, + } +} diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..5d934cd --- /dev/null +++ b/go.mod @@ -0,0 +1,3 @@ +module goflink + +go 1.26.3 diff --git a/main.go b/main.go new file mode 100644 index 0000000..25792c1 --- /dev/null +++ b/main.go @@ -0,0 +1,13 @@ +package main + +import ( + "fmt" + "goflink/core" + "unsafe" +) + +func main() { + + fmt.Printf("GoodRecord Size: %d bytes\n", unsafe.Sizeof(core.StreamRecord[int]{})) + +} -- cgit v1.2.3