mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-13 01:50:40 +02:00
Update fetch.go
Skip long-polling if any requested topic does not exist. Only long-poll when MinBytes > 0, data isn’t available yet, and all topics exist. Cap the long-polling wait to 1s in tests to prevent hanging on shutdown.
This commit is contained in:
@@ -22,6 +22,15 @@ func (h *Handler) handleFetch(correlationID uint32, apiVersion uint16, requestBo
|
||||
|
||||
// Basic long-polling to avoid client busy-looping when there's no data.
|
||||
var throttleTimeMs int32 = 0
|
||||
// Only long-poll when all referenced topics exist; unknown topics should not block
|
||||
allTopicsExist := func() bool {
|
||||
for _, topic := range fetchRequest.Topics {
|
||||
if !h.seaweedMQHandler.TopicExists(topic.Name) {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
hasDataAvailable := func() bool {
|
||||
for _, topic := range fetchRequest.Topics {
|
||||
for _, p := range topic.Partitions {
|
||||
@@ -36,9 +45,15 @@ func (h *Handler) handleFetch(correlationID uint32, apiVersion uint16, requestBo
|
||||
}
|
||||
return false
|
||||
}
|
||||
if fetchRequest.MinBytes > 0 && fetchRequest.MaxWaitTime > 0 && !hasDataAvailable() {
|
||||
// Cap long-polling to avoid blocking connection shutdowns in tests
|
||||
maxWaitMs := fetchRequest.MaxWaitTime
|
||||
if maxWaitMs > 1000 {
|
||||
maxWaitMs = 1000
|
||||
}
|
||||
shouldLongPoll := fetchRequest.MinBytes > 0 && maxWaitMs > 0 && !hasDataAvailable() && allTopicsExist()
|
||||
if shouldLongPoll {
|
||||
start := time.Now()
|
||||
deadline := start.Add(time.Duration(fetchRequest.MaxWaitTime) * time.Millisecond)
|
||||
deadline := start.Add(time.Duration(maxWaitMs) * time.Millisecond)
|
||||
for time.Now().Before(deadline) {
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
if hasDataAvailable() {
|
||||
|
||||
Reference in New Issue
Block a user