feat(go): decouple gateway publish with async sink

This commit is contained in:
lingniu
2026-07-02 14:11:43 +08:00
parent 5a958616a3
commit 5145d156bc
3 changed files with 238 additions and 0 deletions

View File

@@ -178,6 +178,20 @@ func buildSink(ctx context.Context, logger *slog.Logger) (eventbus.Sink, error)
logger.Info("kafka durable spool enabled", "dir", spoolDir, "replay_interval_ms", interval.Milliseconds(), "replay_batch_size", replayBatchSize)
out = durable
}
if envBool("KAFKA_ASYNC_ENABLED", true) {
queueSize := envInt("KAFKA_ASYNC_QUEUE_SIZE", 10000)
workers := envInt("KAFKA_ASYNC_WORKERS", 1)
timeout := time.Duration(envInt("KAFKA_ASYNC_PUBLISH_TIMEOUT_MS", 30000)) * time.Millisecond
out = eventbus.NewAsyncSink(out, eventbus.AsyncConfig{
QueueSize: queueSize,
Workers: workers,
OperationTimeout: timeout,
OnError: func(err error) {
logger.Warn("kafka async publish failed", "error", err)
},
})
logger.Info("kafka async publish enabled", "queue_size", queueSize, "workers", workers, "publish_timeout_ms", timeout.Milliseconds())
}
return out, nil
}

View File

@@ -0,0 +1,119 @@
package eventbus
import (
"context"
"errors"
"sync"
"time"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope"
)
type AsyncConfig struct {
QueueSize int
Workers int
OperationTimeout time.Duration
OnError func(error)
}
type AsyncSink struct {
delegate Sink
jobs chan asyncJob
timeout time.Duration
onError func(error)
closeOnce sync.Once
closed chan struct{}
done chan struct{}
wg sync.WaitGroup
}
type asyncJob struct {
kind string
env envelope.FrameEnvelope
}
var ErrAsyncSinkClosed = errors.New("async sink is closed")
func NewAsyncSink(delegate Sink, cfg AsyncConfig) *AsyncSink {
if delegate == nil {
panic("async delegate sink must not be nil")
}
if cfg.QueueSize <= 0 {
cfg.QueueSize = 10_000
}
if cfg.Workers <= 0 {
cfg.Workers = 1
}
if cfg.OperationTimeout <= 0 {
cfg.OperationTimeout = 30 * time.Second
}
s := &AsyncSink{
delegate: delegate,
jobs: make(chan asyncJob, cfg.QueueSize),
timeout: cfg.OperationTimeout,
onError: cfg.OnError,
closed: make(chan struct{}),
done: make(chan struct{}),
}
s.wg.Add(cfg.Workers)
for i := 0; i < cfg.Workers; i++ {
go s.worker()
}
go func() {
s.wg.Wait()
close(s.done)
}()
return s
}
func (s *AsyncSink) PublishRaw(ctx context.Context, env envelope.FrameEnvelope) error {
return s.enqueue(ctx, asyncJob{kind: "raw", env: env})
}
func (s *AsyncSink) PublishUnified(ctx context.Context, env envelope.FrameEnvelope) error {
return s.enqueue(ctx, asyncJob{kind: "unified", env: env})
}
func (s *AsyncSink) Close() error {
s.closeOnce.Do(func() {
close(s.closed)
close(s.jobs)
})
<-s.done
return s.delegate.Close()
}
func (s *AsyncSink) enqueue(ctx context.Context, job asyncJob) error {
select {
case <-s.closed:
return ErrAsyncSinkClosed
default:
}
select {
case s.jobs <- job:
return nil
case <-s.closed:
return ErrAsyncSinkClosed
case <-ctx.Done():
return ctx.Err()
}
}
func (s *AsyncSink) worker() {
defer s.wg.Done()
for job := range s.jobs {
ctx, cancel := context.WithTimeout(context.Background(), s.timeout)
var err error
switch job.kind {
case "raw":
err = s.delegate.PublishRaw(ctx, job.env)
case "unified":
err = s.delegate.PublishUnified(ctx, job.env)
}
cancel()
if err != nil && s.onError != nil {
s.onError(err)
}
}
}

View File

@@ -0,0 +1,105 @@
package eventbus
import (
"context"
"testing"
"time"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope"
)
func TestAsyncSinkPublishRawReturnsAfterQueueing(t *testing.T) {
delegate := newBlockingSink()
sink := NewAsyncSink(delegate, AsyncConfig{QueueSize: 1, Workers: 1, OperationTimeout: time.Second})
start := time.Now()
if err := sink.PublishRaw(context.Background(), envelope.FrameEnvelope{Protocol: envelope.ProtocolGB32960, VIN: "LNBSCB3D4R1234567"}); err != nil {
t.Fatalf("PublishRaw() error = %v", err)
}
if elapsed := time.Since(start); elapsed > 100*time.Millisecond {
t.Fatalf("PublishRaw() blocked for %s, want quick queueing", elapsed)
}
select {
case <-delegate.rawStarted:
case <-time.After(time.Second):
t.Fatal("delegate raw publish was not started")
}
delegate.release()
if err := sink.Close(); err != nil {
t.Fatalf("Close() error = %v", err)
}
}
func TestAsyncSinkPublishesInFIFOOrderWithSingleWorker(t *testing.T) {
delegate := &orderedSink{}
sink := NewAsyncSink(delegate, AsyncConfig{QueueSize: 4, Workers: 1, OperationTimeout: time.Second})
env := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, Phone: "13307795425"}
if err := sink.PublishRaw(context.Background(), env); err != nil {
t.Fatalf("PublishRaw() error = %v", err)
}
if err := sink.PublishUnified(context.Background(), env); err != nil {
t.Fatalf("PublishUnified() error = %v", err)
}
if err := sink.Close(); err != nil {
t.Fatalf("Close() error = %v", err)
}
if got, want := delegate.calls, []string{"raw", "unified"}; len(got) != len(want) || got[0] != want[0] || got[1] != want[1] {
t.Fatalf("calls = %#v, want %#v", got, want)
}
}
type blockingSink struct {
rawStarted chan struct{}
releaseRaw chan struct{}
}
func newBlockingSink() *blockingSink {
return &blockingSink{
rawStarted: make(chan struct{}),
releaseRaw: make(chan struct{}),
}
}
func (s *blockingSink) PublishRaw(context.Context, envelope.FrameEnvelope) error {
close(s.rawStarted)
<-s.releaseRaw
return nil
}
func (s *blockingSink) PublishUnified(context.Context, envelope.FrameEnvelope) error {
return nil
}
func (s *blockingSink) Close() error {
s.release()
return nil
}
func (s *blockingSink) release() {
select {
case <-s.releaseRaw:
default:
close(s.releaseRaw)
}
}
type orderedSink struct {
calls []string
}
func (s *orderedSink) PublishRaw(context.Context, envelope.FrameEnvelope) error {
s.calls = append(s.calls, "raw")
return nil
}
func (s *orderedSink) PublishUnified(context.Context, envelope.FrameEnvelope) error {
s.calls = append(s.calls, "unified")
return nil
}
func (s *orderedSink) Close() error {
return nil
}