MCPcopy Create free account

hub / github.com/cloudflare/go-stream / functions

Functions520 in github.com/cloudflare/go-stream

MethodExit
()
stream/mapper/generator.go:27
FunctionFatalf
(format string, v ...interface{})
util/slog/slog.go:80
MethodFlushAll
(outch chan<- stream.Object)
cube/operator.go:29
MethodFlushAll
(out chan<- Object)
stream/interfaceBatcher.go:59
MethodFlushItems
()
cube/partitionedcube.go:172
MethodGetAggregates
()
cube/cube.go:86
MethodGetConstraint
(t *Table)
cube/pg/table.go:125
MethodGetCurrentEra
()
cluster/manager.go:20
MethodGetCurrentEra
()
cluster/manager.go:31
MethodGetCurrentEra
()
cluster/manager.go:65
MethodGetDimensions
()
cube/cube.go:82
MethodGetEra
(t time.Time)
cluster/manager.go:19
MethodGetEra
(t time.Time)
cluster/manager.go:27
MethodGetEra
(t time.Time)
cluster/manager.go:58
MethodGetInDepth
()
stream/chain.go:173
MethodGetInDepth
()
stream/util.go:24
MethodGetNode
(posit int)
cluster/era.go:45
MethodGetNodes
()
cluster/era.go:9
MethodGetNodes
()
cluster/era.go:24
MethodGetNodes
()
cluster/era.go:41
MethodGetTableName
(basename string)
cube/pg/table.go:121
MethodGetWorker
()
stream/mapper/generator.go:19
MethodGetWorker
()
stream/mapper/generator.go:37
MethodGetWorker
()
stream/mapper/generator.go:53
MethodGetWorker
()
stream/mapper/generator.go:65
MethodHasItems
()
cube/partitionedcube.go:187
MethodHasItems
()
cube/partitionedcube.go:216
MethodHasItems
()
cube/operator.go:33
MethodHasItems
()
stream/interfaceBatcher.go:63
MethodIn
()
stream/chain.go:168
MethodIn
()
stream/util.go:20
FunctionInit
(logName string, logLevel string, logPrefix string, metrics *util.StreamingMetrics, metricsAddr string, logAd
util/slog/slog.go:45
MethodInsert
(dimensions Dimensions, aggregates Aggregates)
cube/partitionedcube.go:19
MethodInsert
(dimensions Dimensions, aggregates Aggregates)
cube/partitionedcube.go:103
MethodInsert
(dimensions Dimensions, aggregates Aggregates)
cube/partitionedcube.go:146
MethodInsert
(dimensions Dimensions, aggregates Aggregates)
cube/cube.go:58
MethodIp
()
cluster/node.go:16
MethodIp
()
cluster/node.go:34
MethodIp
()
cluster/node.go:63
MethodIsOrdered
()
stream/mapper/orderpreserving.go:46
MethodIsOrdered
()
stream/mapper/operator.go:75
MethodIsParallel
()
stream/operator.go:32
MethodIsParallel
()
stream/mapper/operator.go:71
FunctionJsonGeneralDecoder
* Example Decoder Usage intDecGenFn := func () interface{} { decoder := encoding.JsonGeneralDecoder() return func(in []byte, closenotifier chan<- bo
stream/encoding/json.go:24
FunctionJsonGeneralEncoder
()
stream/encoding/json.go:34
MethodLen
()
transport/client.go:102
MethodLen
()
util/util.go:48
MethodLen
()
util/util.go:88
MethodLen
()
util/util.go:147
MethodMakeOrdered
()
stream/mapper/orderpreserving.go:50
MethodMakeOrdered
()
stream/mapper/operator.go:79
MethodMap
(input stream.Object, out Outputer)
stream/mapper/worker.go:40
MethodMap
(input stream.Object, out Outputer)
stream/mapper/worker.go:103
MethodMerge
(update *PartitionedCube)
cube/partitionedcube.go:38
MethodMerge
(with Aggregate)
cube/aggregate.go:13
MethodMerge
(with Aggregate)
cube/aggregate.go:28
MethodName
()
cube/pg/table.go:34
MethodName
()
cluster/node.go:15
MethodName
()
cluster/node.go:30
MethodName
()
cluster/node.go:59
FunctionNewBaseInOutOp
(slack int)
stream/util.go:58
FunctionNewDistributor
(mapp func(Object) DistribKey, creator func(DistribKey) DistributorChildOp)
stream/distributor.go:25
FunctionNewDropOp
()
stream/util/util.go:10
FunctionNewDynamicBBManager
(bbHosts []string)
cluster/manager.go:123
FunctionNewFanoutOp
()
stream/fanout.go:22
FunctionNewHllAggregate
(val string)
cube/aggregate.go:33
FunctionNewHllAggregateFromBytes
(val []byte)
cube/aggregate.go:42
FunctionNewHllDimension
(i *hll.Hll)
cube/dimension.go:45
FunctionNewIOReaderSource
(reader io.ReadCloser)
stream/source/reader.go:80
FunctionNewIOReaderSourceLengthDelim
(reader io.ReadCloser)
stream/source/reader.go:85
FunctionNewInChainWrapper
(c Chain)
stream/chain.go:164
FunctionNewInterfaceBatchOp
(pn ProcessedNotifier)
stream/interfaceBatcher.go:76
FunctionNewInterfaceBuffer
(size int)
util/util.go:65
FunctionNewInterfaceReaderSource
(reader InterfaceReader)
stream/source/interfacereader.go:35
FunctionNewInterfaceTimingOp
* func NewInterfaceTimingOp() (oper stream.Operator, count *uint32, duration *time.Duration) { var counter = new(uint32) var dur time.Duration var
stream/timing/timing.go:65
FunctionNewInterfaceWriterSink
(writer InterfaceWriter)
stream/sink/interfacewriter.go:35
FunctionNewJsonDecodeRop
(gen interface{})
stream/encoding/json.go:41
FunctionNewJsonEncodeRop
()
stream/encoding/json.go:45
FunctionNewMakeInterfaceOp
()
stream/util/util.go:17
FunctionNewMakeProtobufMessageOp
()
stream/encoding/protobuf.go:53
FunctionNewMemoryBuffer
(size int)
util/util.go:13
FunctionNewMultiPartWriterSink
(writer io.Writer)
stream/sink/writer.go:124
FunctionNewNonBlockingProcessedNotifier
(slack int)
stream/ProcessedNotifier.go:18
FunctionNewOp
(proc interface{}, tn string)
stream/mapper/operator.go:7
FunctionNewOpExitor
(callback interface{}, exitCallback func(), tn string)
stream/mapper/operator.go:15
FunctionNewOpFactory
(proc interface{}, tn string)
stream/mapper/operator.go:23
FunctionNewOpWorkerCloserFactory
(proc interface{}, tn string)
stream/mapper/operator.go:31
FunctionNewOpWorkerFinalItemsFactory
(proc interface{}, tn string)
stream/mapper/operator.go:39
FunctionNewOrderedOp
(proc interface{}, tn string)
stream/mapper/orderpreserving.go:9
FunctionNewPgBatchOperator
(parse func(stream.Object) (Dimensions, Aggregates), downstreamProcessed stream.ProcessedNotifier)
cube/operator.go:37
FunctionNewProcessedNotifier
()
stream/ProcessedNotifier.go:14
FunctionNewProtobufDecodeOp
(gen interface{})
stream/encoding/protobuf.go:33
FunctionNewProtobufEncodeOp
()
stream/encoding/protobuf.go:37
FunctionNewSequentialBufferChanImpl
(maxItems int)
util/util.go:116
FunctionNewSimpleEra
()
cluster/era.go:16
FunctionNewSimpleNode
(name string, ip string, port string)
cluster/node.go:26
FunctionNewSnappyDecodeOp
()
stream/compress/snappy.go:24
FunctionNewSnappyEncodeOp
()
stream/compress/snappy.go:10
FunctionNewStaticManager
(e Era)
cluster/manager.go:35
FunctionNewStreamingMetrics
(mReg metrics.Registry)
util/util.go:198
← previousnext →301–400 of 520, ranked by callers