Browse Source

fix locking

pull/7329/head
chrislu 4 weeks ago
parent
commit
98b536480d
  1. 4
      weed/mq/kafka/integration/broker_client_subscribe.go

4
weed/mq/kafka/integration/broker_client_subscribe.go

@ -435,8 +435,8 @@ func (bc *BrokerClient) ReadRecordsFromOffset(ctx context.Context, session *Brok
// - Exact match (requestedOffset == session.StartOffset)
// - Reading ahead (requestedOffset > session.StartOffset, e.g., from cache)
glog.V(2).Infof("[FETCH] Using persistent session: requested=%d session=%d (persistent connection)",
requestedOffset, session.StartOffset)
session.mu.Unlock()
requestedOffset, currentStartOffset)
// Note: session.mu was already unlocked at line 294 after reading currentStartOffset
return bc.ReadRecords(ctx, session, maxRecords)
}

Loading…
Cancel
Save