diff options
| author | SangTran-127 <tranquangsang12.7@gmail.com> | 2026-08-08 11:31:29 +0700 |
|---|---|---|
| committer | SangTran-127 <tranquangsang12.7@gmail.com> | 2026-08-08 11:31:29 +0700 |
| commit | 77a3b3e86b06504bdf37e35d83a192806a95a63e (patch) | |
| tree | d37f9f6cfe04064823803f28e2b0cc4a00e3127c /core/record.go | |
| parent | 967c0348c4f5c44801da90dbd9cdaeca09560a81 (diff) | |
| download | goflink-77a3b3e86b06504bdf37e35d83a192806a95a63e.tar.gz goflink-77a3b3e86b06504bdf37e35d83a192806a95a63e.zip | |
initialize goflink module with core stream record implementation
Diffstat (limited to 'core/record.go')
| -rw-r--r-- | core/record.go | 40 |
1 files changed, 40 insertions, 0 deletions
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, + } +} |