mirror of
https://gitee.com/milvus-io/milvus.git
synced 2025-12-08 01:58:34 +08:00
issue: #43897, #44123 pr: #45224 also pick pr: #45216,#45154,#45033,#45145,#45092,#45058,#45029 enhance: Close channel replicator more gracefully (#45029) issue: https://github.com/milvus-io/milvus/issues/44123 enhance: Show create time for import job (#45058) issue: https://github.com/milvus-io/milvus/issues/45056 fix: wal state may be unconsistent after recovering from crash (#45092) issue: #45088, #45086 - Message on control channel should trigger the checkpoint update. - LastConfrimedMessageID should be recovered from the minimum of checkpoint or the LastConfirmedMessageID of uncommitted txn. - Add more log info for wal debugging. fix: make ack of broadcaster cannot canceled by client (#45145) issue: #45141 - make ack of broadcaster cannot canceled by rpc. - make clone for assignment snapshot of wal balancer. - add server id for GetReplicateCheckpoint to avoid failure. enhance: support collection and index with WAL-based DDL framework (#45033) issue: #43897 - Part of collection/index related DDL is implemented by WAL-based DDL framework now. - Support following message type in wal, CreateCollection, DropCollection, CreatePartition, DropPartition, CreateIndex, AlterIndex, DropIndex. - Part of collection/index related DDL can be synced by new CDC now. - Refactor some UT for collection/index DDL. - Add Tombstone scheduler to manage the tombstone GC for collection or partition meta. - Move the vchannel allocation into streaming pchannel manager. enhance: support load/release collection/partition with WAL-based DDL framework (#45154) issue: #43897 - Load/Release collection/partition is implemented by WAL-based DDL framework now. - Support AlterLoadConfig/DropLoadConfig in wal now. - Load/Release operation can be synced by new CDC now. - Refactor some UT for load/release DDL. enhance: Don't start cdc by default (#45216) issue: https://github.com/milvus-io/milvus/issues/44123 fix: unrecoverable when replicate from old (#45224) issue: #44962 --------- Signed-off-by: bigsheeper <yihao.dai@zilliz.com> Signed-off-by: chyezh <chyezh@outlook.com> Co-authored-by: yihao.dai <yihao.dai@zilliz.com>
76 lines
2.6 KiB
Go
76 lines
2.6 KiB
Go
package service
|
|
|
|
import (
|
|
"context"
|
|
|
|
"github.com/milvus-io/milvus/internal/streamingcoord/server/broadcaster/broadcast"
|
|
"github.com/milvus-io/milvus/pkg/v2/proto/streamingpb"
|
|
"github.com/milvus-io/milvus/pkg/v2/streaming/util/message"
|
|
)
|
|
|
|
// BroadcastService is the interface of the broadcast service.
|
|
type BroadcastService interface {
|
|
streamingpb.StreamingCoordBroadcastServiceServer
|
|
}
|
|
|
|
// NewBroadcastService creates a new broadcast service.
|
|
func NewBroadcastService() BroadcastService {
|
|
return &broadcastServceImpl{}
|
|
}
|
|
|
|
// broadcastServiceeeeImpl is the implementation of the broadcast service.
|
|
type broadcastServceImpl struct{}
|
|
|
|
// Broadcast broadcasts the message to all channels.
|
|
func (s *broadcastServceImpl) Broadcast(ctx context.Context, req *streamingpb.BroadcastRequest) (*streamingpb.BroadcastResponse, error) {
|
|
msg := message.NewBroadcastMutableMessageBeforeAppend(req.Message.Payload, req.Message.Properties)
|
|
api, err := broadcast.StartBroadcastWithResourceKeys(ctx, msg.BroadcastHeader().ResourceKeys.Collect()...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer api.Close()
|
|
|
|
results, err := api.Broadcast(ctx, msg)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
protoResult := make(map[string]*streamingpb.ProduceMessageResponseResult, len(results.AppendResults))
|
|
for vchannel, result := range results.AppendResults {
|
|
protoResult[vchannel] = &streamingpb.ProduceMessageResponseResult{
|
|
Id: result.MessageID.IntoProto(),
|
|
Timetick: result.TimeTick,
|
|
LastConfirmedId: result.LastConfirmedMessageID.IntoProto(),
|
|
}
|
|
}
|
|
return &streamingpb.BroadcastResponse{
|
|
BroadcastId: results.BroadcastID,
|
|
Results: protoResult,
|
|
}, nil
|
|
}
|
|
|
|
// Ack acknowledges the message at the specified vchannel.
|
|
func (s *broadcastServceImpl) Ack(ctx context.Context, req *streamingpb.BroadcastAckRequest) (*streamingpb.BroadcastAckResponse, error) {
|
|
broadcaster, err := broadcast.GetWithContext(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
// Once the ack is reached at streamingcoord, the ack operation should not be cancelable.
|
|
ctx = context.WithoutCancel(ctx)
|
|
if req.Message == nil {
|
|
// before 2.6.1, the request don't have the message field, only have the broadcast id and vchannel.
|
|
// so we need to use the legacy ack interface.
|
|
if err := broadcaster.LegacyAck(ctx, req.BroadcastId, req.Vchannel); err != nil {
|
|
return nil, err
|
|
}
|
|
return &streamingpb.BroadcastAckResponse{}, nil
|
|
}
|
|
if err := broadcaster.Ack(ctx, message.NewImmutableMesasge(
|
|
message.MustUnmarshalMessageID(req.Message.Id),
|
|
req.Message.Payload,
|
|
req.Message.Properties,
|
|
)); err != nil {
|
|
return nil, err
|
|
}
|
|
return &streamingpb.BroadcastAckResponse{}, nil
|
|
}
|