mirror of
https://github.com/milvus-io/milvus.git
synced 2026-07-21 02:05:41 +00:00
issue: #46540 Empty timetick is just used to sync up the time clock between different component in milvus. So empty timetick can be ignored if we achieve the lsn/mvcc semantic for timetick. Currently, some components need the empty timetick to trigger some operation, such as flush/tsafe. So we only slow down the empty time tick for 5 seconds. <!-- This is an auto-generated comment: release notes by coderabbit.ai --> - Core invariant: with LSN/MVCC semantics consumers only need (a) the first timetick that advances the latest-required-MVCC to unblock MVCC-dependent waits and (b) occasional periodic timeticks (~≤5s) for clock synchronization—therefore frequent non-persisted empty timeticks can be suppressed without breaking MVCC correctness. - Logic removed/simplified: per-message dispatch/consumption of frequent non-persisted empty timeticks is suppressed — an MVCC-aware filter emptyTimeTickSlowdowner (internal/util/pipeline/consuming_slowdown.go) short-circuits frequent empty timeticks in the stream pipeline (internal/util/pipeline/stream_pipeline.go), and the WAL flusher rate-limits non-persisted timetick dispatch to one emission per ~5s (internal/streamingnode/server/flusher/flusherimpl/wal_flusher.go); the delegator exposes GetLatestRequiredMVCCTimeTick to drive the filter (internal/querynodev2/delegator/delegator.go). - Why this does NOT introduce data loss or regressions: the slowdowner always refreshes latestRequiredMVCCTimeTick via GetLatestRequiredMVCCTimeTick and (1) never filters timeticks < latestRequiredMVCCTimeTick (so existing tsafe/flush waits stay unblocked) and (2) always lets the first timetick ≥ latestRequiredMVCCTimeTick pass to notify pending MVCC waits; separately, WAL flusher suppression applies only to non-persisted timeticks and still emits when the 5s threshold elapses, preserving periodic clock-sync messages used by flush/tsafe. - Enhancement summary (where it takes effect): adds GetLatestRequiredMVCCTimeTick on ShardDelegator and LastestMVCCTimeTickGetter, wires emptyTimeTickSlowdowner into NewPipelineWithStream (internal/util/pipeline), and adds WAL flusher rate-limiting + metrics (internal/streamingnode/server/flusher/flusherimpl/wal_flusher.go, pkg/metrics) to reduce CPU/dispatch overhead while keeping MVCC correctness and periodic synchronization. <!-- end of auto-generated comment: release notes by coderabbit.ai --> --------- Signed-off-by: chyezh <chyezh@outlook.com>
123 lines
3.9 KiB
YAML
123 lines
3.9 KiB
YAML
quiet: False
|
|
with-expecter: True
|
|
filename: "mock_{{.InterfaceName}}.go"
|
|
dir: 'internal/mocks/{{trimPrefix .PackagePath "github.com/milvus-io/milvus/internal" | dir }}/mock_{{.PackageName}}'
|
|
mockname: "Mock{{.InterfaceName}}"
|
|
outpkg: "mock_{{.PackageName}}"
|
|
resolve-type-alias: False
|
|
packages:
|
|
github.com/milvus-io/milvus/internal/distributed/streaming:
|
|
interfaces:
|
|
WALAccesser:
|
|
ReplicateService:
|
|
Utility:
|
|
Broadcast:
|
|
Local:
|
|
Scanner:
|
|
github.com/milvus-io/milvus/internal/rootcoord/tombstone:
|
|
interfaces:
|
|
TombstoneSweeper:
|
|
github.com/milvus-io/milvus/internal/streamingcoord/server/balancer:
|
|
interfaces:
|
|
Balancer:
|
|
github.com/milvus-io/milvus/internal/streamingcoord/server/broadcaster:
|
|
interfaces:
|
|
Broadcaster:
|
|
BroadcastAPI:
|
|
AppendOperator:
|
|
Watcher:
|
|
github.com/milvus-io/milvus/internal/streamingcoord/client:
|
|
interfaces:
|
|
Client:
|
|
BroadcastService:
|
|
AssignmentService:
|
|
github.com/milvus-io/milvus/internal/streamingcoord/client/broadcast:
|
|
interfaces:
|
|
Watcher:
|
|
github.com/milvus-io/milvus/internal/streamingnode/client/manager:
|
|
interfaces:
|
|
ManagerClient:
|
|
github.com/milvus-io/milvus/internal/streamingnode/client/handler:
|
|
interfaces:
|
|
HandlerClient:
|
|
github.com/milvus-io/milvus/internal/streamingnode/client/handler/assignment:
|
|
interfaces:
|
|
Watcher:
|
|
github.com/milvus-io/milvus/internal/streamingnode/client/handler/producer:
|
|
interfaces:
|
|
Producer:
|
|
github.com/milvus-io/milvus/internal/streamingnode/client/handler/consumer:
|
|
interfaces:
|
|
Consumer:
|
|
github.com/milvus-io/milvus/internal/streamingnode/server/wal:
|
|
interfaces:
|
|
OpenerBuilder:
|
|
Opener:
|
|
Scanner:
|
|
WAL:
|
|
github.com/milvus-io/milvus/internal/streamingnode/server/wal/interceptors:
|
|
interfaces:
|
|
Interceptor:
|
|
InterceptorWithReady:
|
|
InterceptorWithMetrics:
|
|
InterceptorBuilder:
|
|
github.com/milvus-io/milvus/internal/streamingnode/server/wal/interceptors/replicate/replicates:
|
|
interfaces:
|
|
ReplicatesManager:
|
|
ReplicateAcker:
|
|
github.com/milvus-io/milvus/internal/streamingnode/server/wal/interceptors/shard/shards:
|
|
interfaces:
|
|
ShardManager:
|
|
github.com/milvus-io/milvus/internal/streamingnode/server/wal/interceptors/shard/utils:
|
|
interfaces:
|
|
SealOperator:
|
|
github.com/milvus-io/milvus/internal/streamingnode/server/wal/interceptors/timetick/inspector:
|
|
interfaces:
|
|
TimeTickSyncOperator:
|
|
github.com/milvus-io/milvus/internal/streamingnode/server/wal/recovery:
|
|
interfaces:
|
|
RecoveryStorage:
|
|
github.com/milvus-io/milvus/internal/streamingnode/server/wal/interceptors/wab:
|
|
interfaces:
|
|
ROWriteAheadBuffer:
|
|
google.golang.org/grpc:
|
|
interfaces:
|
|
ClientStream:
|
|
github.com/milvus-io/milvus/internal/streamingnode/server/walmanager:
|
|
interfaces:
|
|
Manager:
|
|
github.com/milvus-io/milvus/internal/metastore:
|
|
interfaces:
|
|
ReplicationCatalog:
|
|
StreamingCoordCataLog:
|
|
StreamingNodeCataLog:
|
|
github.com/milvus-io/milvus/internal/util/segcore:
|
|
interfaces:
|
|
CSegment:
|
|
github.com/milvus-io/milvus/internal/storage:
|
|
interfaces:
|
|
ChunkManager:
|
|
github.com/milvus-io/milvus/internal/util/streamingutil/service/discoverer:
|
|
interfaces:
|
|
Discoverer:
|
|
AssignmentDiscoverWatcher:
|
|
github.com/milvus-io/milvus/internal/util/streamingutil/service/lazygrpc:
|
|
interfaces:
|
|
Service:
|
|
github.com/milvus-io/milvus/internal/util/streamingutil/service/resolver:
|
|
interfaces:
|
|
Resolver:
|
|
Builder:
|
|
github.com/milvus-io/milvus/internal/util/searchutil/optimizers:
|
|
interfaces:
|
|
QueryHook:
|
|
github.com/milvus-io/milvus/internal/util/pipeline:
|
|
interfaces:
|
|
LastestMVCCTimeTickGetter:
|
|
github.com/milvus-io/milvus/internal/flushcommon/util:
|
|
interfaces:
|
|
MsgHandler:
|
|
google.golang.org/grpc/resolver:
|
|
interfaces:
|
|
ClientConn:
|