-
Notifications
You must be signed in to change notification settings - Fork 131
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
GRPC Streaming Batching #1633
GRPC Streaming Batching #1633
Changes from 3 commits
e413e5c
7974fe1
edb122c
3aeb36b
e23f3cb
536dffe
7d40694
434e4d3
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -21,8 +21,11 @@ type Flags struct { | |
GrpcEnable bool | ||
|
||
// Grpc Streaming | ||
GrpcStreamingEnabled bool | ||
VEOracleEnabled bool // Slinky Vote Extensions | ||
GrpcStreamingEnabled bool | ||
GrpcStreamingFlushIntervalMs uint32 | ||
GrpcStreamingMaxBufferSize uint32 | ||
|
||
VEOracleEnabled bool // Slinky Vote Extensions | ||
} | ||
|
||
// List of CLI flags. | ||
|
@@ -37,7 +40,9 @@ const ( | |
GrpcEnable = "grpc.enable" | ||
|
||
// Grpc Streaming | ||
GrpcStreamingEnabled = "grpc-streaming-enabled" | ||
GrpcStreamingEnabled = "grpc-streaming-enabled" | ||
GrpcStreamingFlushIntervalMs = "grpc-streaming-flush-interval-ms" | ||
GrpcStreamingMaxBufferSize = "grpc-streaming-max-buffer-size" | ||
|
||
// Slinky VEs enabled | ||
VEOracleEnabled = "slinky-vote-extension-oracle-enabled" | ||
|
@@ -50,8 +55,11 @@ const ( | |
DefaultNonValidatingFullNode = false | ||
DefaultDdErrorTrackingFormat = false | ||
|
||
DefaultGrpcStreamingEnabled = false | ||
DefaultVEOracleEnabled = true | ||
DefaultGrpcStreamingEnabled = false | ||
DefaultGrpcStreamingFlushIntervalMs = 50 | ||
DefaultGrpcStreamingMaxBufferSize = 10000 | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. is 10k reasonable here? should we go higher in case of spikes in traffic? There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. i guess if we are doing 10ms, 10k should be enough There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. i'll use these for now and then tune later based on metrics |
||
|
||
DefaultVEOracleEnabled = true | ||
) | ||
|
||
// AddFlagsToCmd adds flags to app initialization. | ||
|
@@ -85,6 +93,16 @@ func AddFlagsToCmd(cmd *cobra.Command) { | |
DefaultGrpcStreamingEnabled, | ||
"Whether to enable grpc streaming for full nodes", | ||
) | ||
cmd.Flags().Uint32( | ||
GrpcStreamingFlushIntervalMs, | ||
DefaultGrpcStreamingFlushIntervalMs, | ||
"Flush interval (in ms) for grpc streaming", | ||
) | ||
cmd.Flags().Uint32( | ||
GrpcStreamingMaxBufferSize, | ||
DefaultGrpcStreamingMaxBufferSize, | ||
"Maximum buffer size before grpc streaming cancels all connections", | ||
) | ||
cmd.Flags().Bool( | ||
VEOracleEnabled, | ||
DefaultVEOracleEnabled, | ||
|
@@ -104,6 +122,12 @@ func (f *Flags) Validate() error { | |
if !f.GrpcEnable { | ||
return fmt.Errorf("grpc.enable must be set to true - grpc streaming requires gRPC server") | ||
} | ||
if f.GrpcStreamingMaxBufferSize == 0 { | ||
return fmt.Errorf("grpc streaming buffer size must be positive number") | ||
} | ||
if f.GrpcStreamingFlushIntervalMs == 0 { | ||
return fmt.Errorf("grpc streaming flush interval must be positive number") | ||
} | ||
} | ||
return nil | ||
} | ||
|
@@ -124,8 +148,11 @@ func GetFlagValuesFromOptions( | |
GrpcAddress: config.DefaultGRPCAddress, | ||
GrpcEnable: true, | ||
|
||
GrpcStreamingEnabled: DefaultGrpcStreamingEnabled, | ||
VEOracleEnabled: true, | ||
GrpcStreamingEnabled: DefaultGrpcStreamingEnabled, | ||
GrpcStreamingFlushIntervalMs: DefaultGrpcStreamingFlushIntervalMs, | ||
GrpcStreamingMaxBufferSize: DefaultGrpcStreamingMaxBufferSize, | ||
|
||
VEOracleEnabled: true, | ||
} | ||
|
||
// Populate the flags if they exist. | ||
|
@@ -171,6 +198,18 @@ func GetFlagValuesFromOptions( | |
} | ||
} | ||
|
||
if option := appOpts.Get(GrpcStreamingFlushIntervalMs); option != nil { | ||
if v, err := cast.ToUint32E(option); err == nil { | ||
result.GrpcStreamingFlushIntervalMs = v | ||
} | ||
} | ||
|
||
if option := appOpts.Get(GrpcStreamingMaxBufferSize); option != nil { | ||
if v, err := cast.ToUint32E(option); err == nil { | ||
result.GrpcStreamingMaxBufferSize = v | ||
} | ||
} | ||
|
||
if option := appOpts.Get(VEOracleEnabled); option != nil { | ||
if v, err := cast.ToBoolE(option); err == nil { | ||
result.VEOracleEnabled = v | ||
|
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -25,6 +25,15 @@ type GrpcStreamingManagerImpl struct { | |
// orderbookSubscriptions maps subscription IDs to their respective orderbook subscriptions. | ||
orderbookSubscriptions map[uint32]*OrderbookSubscription | ||
nextSubscriptionId uint32 | ||
|
||
// grpc stream will batch and flush out messages every 10 ms. | ||
ticker *time.Ticker | ||
done chan bool | ||
// map of clob pair id to stream updates. | ||
streamUpdateCache map[uint32][]clobtypes.StreamUpdate | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. not related to this PR, but in memclob we should probably only generate updates for clob pair ids that are relevant. e.g. those with at lease one subscriber |
||
numUpdatesInCache uint32 | ||
|
||
maxUpdatesInCache uint32 | ||
} | ||
|
||
// OrderbookSubscription represents a active subscription to the orderbook updates stream. | ||
|
@@ -41,12 +50,36 @@ type OrderbookSubscription struct { | |
|
||
func NewGrpcStreamingManager( | ||
logger log.Logger, | ||
flushIntervalMs uint32, | ||
maxUpdatesInCache uint32, | ||
) *GrpcStreamingManagerImpl { | ||
logger = logger.With(log.ModuleKey, "grpc-streaming") | ||
return &GrpcStreamingManagerImpl{ | ||
grpcStreamingManager := &GrpcStreamingManagerImpl{ | ||
logger: logger, | ||
orderbookSubscriptions: make(map[uint32]*OrderbookSubscription), | ||
nextSubscriptionId: 0, | ||
|
||
ticker: time.NewTicker(time.Duration(flushIntervalMs) * time.Millisecond), | ||
done: make(chan bool), | ||
streamUpdateCache: make(map[uint32][]clobtypes.StreamUpdate), | ||
numUpdatesInCache: 0, | ||
|
||
maxUpdatesInCache: maxUpdatesInCache, | ||
} | ||
|
||
// Start the goroutine for pushing order updates through | ||
go func() { | ||
for { | ||
select { | ||
case <-grpcStreamingManager.ticker.C: | ||
grpcStreamingManager.FlushStreamUpdates() | ||
case <-grpcStreamingManager.done: | ||
return | ||
} | ||
} | ||
}() | ||
|
||
return grpcStreamingManager | ||
Comment on lines
+54
to
+83
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The constructor |
||
} | ||
|
||
func (sm *GrpcStreamingManagerImpl) Enabled() bool { | ||
|
@@ -88,6 +121,10 @@ func (sm *GrpcStreamingManagerImpl) Subscribe( | |
return nil | ||
} | ||
|
||
func (sm *GrpcStreamingManagerImpl) Stop() { | ||
sm.done <- true | ||
} | ||
|
||
// SendOrderbookUpdates groups updates by their clob pair ids and | ||
// sends messages to the subscribers. | ||
func (sm *GrpcStreamingManagerImpl) SendOrderbookUpdates( | ||
|
@@ -133,9 +170,7 @@ func (sm *GrpcStreamingManagerImpl) SendOrderbookUpdates( | |
} | ||
} | ||
|
||
sm.sendStreamUpdate( | ||
updatesByClobPairId, | ||
) | ||
sm.AddUpdatesToCache(updatesByClobPairId, uint32(len(updates))) | ||
} | ||
|
||
// SendOrderbookFillUpdates groups fills by their clob pair ids and | ||
|
@@ -172,29 +207,49 @@ func (sm *GrpcStreamingManagerImpl) SendOrderbookFillUpdates( | |
updatesByClobPairId[clobPairId] = append(updatesByClobPairId[clobPairId], streamUpdate) | ||
} | ||
|
||
sm.sendStreamUpdate( | ||
updatesByClobPairId, | ||
) | ||
sm.AddUpdatesToCache(updatesByClobPairId, uint32(len(orderbookFills))) | ||
} | ||
|
||
// sendStreamUpdate takes in a map of clob pair id to stream updates and emits them to subscribers. | ||
func (sm *GrpcStreamingManagerImpl) sendStreamUpdate( | ||
func (sm *GrpcStreamingManagerImpl) AddUpdatesToCache( | ||
updatesByClobPairId map[uint32][]clobtypes.StreamUpdate, | ||
numUpdatesToAdd uint32, | ||
) { | ||
sm.Lock() | ||
defer sm.Unlock() | ||
|
||
for clobPairId, streamUpdates := range updatesByClobPairId { | ||
sm.streamUpdateCache[clobPairId] = append(sm.streamUpdateCache[clobPairId], streamUpdates...) | ||
} | ||
sm.numUpdatesInCache += numUpdatesToAdd | ||
|
||
// Remove all subscriptions and wipe the buffer if buffer overflows. | ||
if sm.numUpdatesInCache > sm.maxUpdatesInCache { | ||
sm.logger.Error("GRPC Streaming buffer full capacity. Dropping messages and all subscriptions. " + | ||
"Disconnect all clients and increase buffer size via the grpc-stream-buffer-size flag.") | ||
for id := range sm.orderbookSubscriptions { | ||
delete(sm.orderbookSubscriptions, id) | ||
} | ||
clear(sm.streamUpdateCache) | ||
sm.numUpdatesInCache = 0 | ||
} | ||
} | ||
|
||
// FlushStreamUpdates takes in a map of clob pair id to stream updates and emits them to subscribers. | ||
func (sm *GrpcStreamingManagerImpl) FlushStreamUpdates() { | ||
sm.Lock() | ||
defer sm.Unlock() | ||
|
||
metrics.IncrCounter( | ||
metrics.GrpcEmitProtocolUpdateCount, | ||
1, | ||
) | ||
|
||
sm.Lock() | ||
defer sm.Unlock() | ||
|
||
// Send updates to subscribers. | ||
idsToRemove := make([]uint32, 0) | ||
for id, subscription := range sm.orderbookSubscriptions { | ||
streamUpdatesForSubscription := make([]clobtypes.StreamUpdate, 0) | ||
for _, clobPairId := range subscription.clobPairIds { | ||
if update, ok := updatesByClobPairId[clobPairId]; ok { | ||
if update, ok := sm.streamUpdateCache[clobPairId]; ok { | ||
streamUpdatesForSubscription = append(streamUpdatesForSubscription, update...) | ||
} | ||
} | ||
|
@@ -209,6 +264,7 @@ func (sm *GrpcStreamingManagerImpl) sendStreamUpdate( | |
Updates: streamUpdatesForSubscription, | ||
}, | ||
); err != nil { | ||
sm.logger.Error("Error sending out update", "err", err) | ||
idsToRemove = append(idsToRemove, id) | ||
} | ||
} | ||
|
@@ -219,6 +275,10 @@ func (sm *GrpcStreamingManagerImpl) sendStreamUpdate( | |
for _, id := range idsToRemove { | ||
delete(sm.orderbookSubscriptions, id) | ||
} | ||
|
||
clear(sm.streamUpdateCache) | ||
sm.numUpdatesInCache = 0 | ||
|
||
sm.EmitMetrics() | ||
} | ||
|
||
|
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Consider adding error handling for GRPC streaming manager initialization.
Committable suggestion