Documentation
¶
Index ¶
Constants ¶
This section is empty.
Variables ¶
var (
DROP = fmt.Sprintf("%U__DROP__", '\\') // U+005C__DROP__
)
Functions ¶
func DefaultOptions ¶
func DefaultOptions() *options
func NewServer ¶
func NewServer(ms MapStreamer, inputOptions ...Option) numaflow.Server
NewServer creates a new map streaming server.
Types ¶
type MapStreamer ¶
type MapStreamer interface { // MapStream is the function to process each coming message and streams // the result back using a channel. MapStream(ctx context.Context, keys []string, datum Datum, messageCh chan<- Message) }
MapStreamer is the interface of map stream function implementation.
type MapStreamerFunc ¶
type MapStreamerFunc func(ctx context.Context, keys []string, datum Datum, messageCh chan<- Message)
MapStreamerFunc is a utility type used to convert a function to a MapStreamer.
type Message ¶
type Message struct {
// contains filtered or unexported fields
}
Message is used to wrap the data return by MapStream functions
type Messages ¶
type Messages []Message
func MessagesBuilder ¶
func MessagesBuilder() Messages
MessagesBuilder returns an empty instance of Messages
type Option ¶
type Option func(*options)
Option is the interface to apply options.
func WithMaxMessageSize ¶
WithMaxMessageSize sets the server max receive message size and the server max send message size to the given size.
func WithServerInfoFilePath ¶
WithServerInfoFilePath sets the server info file path to the given path.
func WithSockAddr ¶
WithSockAddr start the server with the given sock addr. This is mainly used for testing purposes.
type Service ¶
type Service struct { mapstreampb.UnimplementedMapStreamServer MapperStream MapStreamer }
Service implements the proto gen server interface and contains the map streaming function.
func (*Service) IsReady ¶
func (fs *Service) IsReady(context.Context, *emptypb.Empty) (*mapstreampb.ReadyResponse, error)
IsReady returns true to indicate the gRPC connection is ready.
func (*Service) MapStreamFn ¶
func (fs *Service) MapStreamFn(d *mapstreampb.MapStreamRequest, stream mapstreampb.MapStream_MapStreamFnServer) error
MapStreamFn applies a function to each request element and streams the results back.