package state import ( "context" "goflink/core" ) // TODO: currently using Golang Map for go through the DAG pipline // This should be use index map index map[uint64]uint32 Hash(KeyGroup + StateName + Key) -> Offset // with 1GB of []byte allocation // This should be optimize with ZERO GC Scanning, Memory Alignment type MemoryStateBackend struct { maxParallelism int currentKey string currentGroupKey int // 3D structure StateName -> KeyGroup -> Key -> Value states map[string]map[int]map[string]any } func NewMemoryStateBackend(maxParallelism int) *MemoryStateBackend { return &MemoryStateBackend{ maxParallelism: maxParallelism, states: make(map[string]map[int]map[string]any), } } func (ms *MemoryStateBackend) GetCurrentKey() string { return ms.currentKey } func (ms *MemoryStateBackend) SetCurrentKey(key string) error { ms.currentKey = key ms.currentGroupKey = core.AssignToKeyGroup(key, ms.maxParallelism) return nil } func (ms *MemoryStateBackend) GetValueState(ctx context.Context, key string) (core.ValueState[any], error) { // Assign if state not existed if _, ok := ms.states[key]; !ok { ms.states[key] = make(map[int]map[string]any) } return &memoryStateValue{backend: ms, stateName: key}, nil }