diff options
| author | chrislu <chris.lu@gmail.com> | 2022-05-30 22:47:29 -0700 |
|---|---|---|
| committer | chrislu <chris.lu@gmail.com> | 2022-05-30 22:47:29 -0700 |
| commit | f4a6da6cb276f1891b01097670b044fd4ee6139d (patch) | |
| tree | 8b92c13091e1d9973ce21e4f0c9ac143bd667cfd /weed/messaging | |
| parent | 596c3860cac83a75ae9ce728c8a043133c03d098 (diff) | |
| parent | ca01ce05249c336ed380d9f77efbee68213b8a37 (diff) | |
| download | seaweedfs-f4a6da6cb276f1891b01097670b044fd4ee6139d.tar.xz seaweedfs-f4a6da6cb276f1891b01097670b044fd4ee6139d.zip | |
Merge branch 'master' of https://github.com/chrislusf/seaweedfs
Diffstat (limited to 'weed/messaging')
| -rw-r--r-- | weed/messaging/broker/broker_grpc_server_subscribe.go | 4 |
1 files changed, 2 insertions, 2 deletions
diff --git a/weed/messaging/broker/broker_grpc_server_subscribe.go b/weed/messaging/broker/broker_grpc_server_subscribe.go index f29121c76..20d529239 100644 --- a/weed/messaging/broker/broker_grpc_server_subscribe.go +++ b/weed/messaging/broker/broker_grpc_server_subscribe.go @@ -117,7 +117,7 @@ func (broker *MessageBroker) Subscribe(stream messaging_pb.SeaweedMessaging_Subs lastReadTime = time.Unix(0, processedTsNs) } - lastReadTime, err = lock.logBuffer.LoopProcessLogData("broker", lastReadTime, func() bool { + lastReadTime, _, err = lock.logBuffer.LoopProcessLogData("broker", lastReadTime, 0, func() bool { lock.Mutex.Lock() lock.cond.Wait() lock.Mutex.Unlock() @@ -164,7 +164,7 @@ func (broker *MessageBroker) readPersistedLogBuffer(tp *TopicPartition, startTim // println("partition", tp.Partition, "processing", dayDir, "/", hourMinuteEntry.Name) chunkedFileReader := filer.NewChunkStreamReader(broker, hourMinuteEntry.Chunks) defer chunkedFileReader.Close() - if _, err := filer.ReadEachLogEntry(chunkedFileReader, sizeBuf, startTsNs, eachLogEntryFn); err != nil { + if _, err := filer.ReadEachLogEntry(chunkedFileReader, sizeBuf, startTsNs, 0, eachLogEntryFn); err != nil { chunkedFileReader.Close() if err == io.EOF { return err |
