MCPcopy Create free account

hub / github.com/JasonThon/lightflus / functions

Functions588 in github.com/JasonThon/lightflus

↓ 2 callersFunctioneval
( scope: &mut v8::HandleScope<'s>, code: &str, )
src/stream/src/v8_runtime.rs:498
↓ 2 callersMethodexecute
()
typescript-api/src/stream/dataflow.ts:42
↓ 2 callersFunctionexecution_id_unprovided
()
src/lightflus-core/src/errors/mod.rs:19
↓ 2 callersFunctionextract_arguments
( extractors: &[String], event: &LocalEvent, fn_name: &str, )
src/stream/src/connector.rs:662
↓ 2 callersMethodflatMap
(callbackFn: (value: T) => U[])
typescript-api/src/stream/dataflow.ts:22
↓ 2 callersMethodfrom
(conn_opts: mysql_desc::ConnectionOpts)
src/common/src/db.rs:166
↓ 2 callersFunctionfrom_reader
(reader: R)
src/common/src/utils.rs:148
↓ 2 callersFunctionfrom_str
(value: &str)
src/common/src/utils.rs:157
↓ 2 callersMethodgetId
()
typescript-api/src/stream/context.ts:35
↓ 2 callersMethodget_dataflow
( &self, job_id: &ResourceId, )
src/lightflus-core/src/coordinator/coord.rs:108
↓ 2 callersMethodget_downstream_id_iter
(&self)
src/stream/src/task.rs:77
↓ 2 callersMethodget_gateway
(&self)
src/common/src/net/cluster.rs:67
↓ 2 callersMethodget_id
(&self)
src/common/src/net/cluster.rs:72
↓ 2 callersMethodget_kafka_group
(&self)
src/proto/src/common_impl.rs:195
↓ 2 callersMethodget_key
(&self)
src/proto/src/common_impl.rs:387
↓ 2 callersMethodget_keyed_state
(&self, key: &[u8])
src/stream/src/state.rs:46
↓ 2 callersFunctionget_operator_state_key
(operator_id: NodeIdx, operator: &str, reference: &[u8])
src/stream/src/dataflow.rs:400
↓ 2 callersMethodget_sub_dataflow
( &self, request: RpcRequest<ResourceId>, )
src/lightflus-core/src/taskmanager/rpc.rs:191
↓ 2 callersMethodget_subdataflow_id
(&self)
src/proto/src/common_impl.rs:521
↓ 2 callersMethodgroup
(value: string)
typescript-api/src/connectors/connectors.ts:41
↓ 2 callersMethodhas_source
(&self)
src/proto/src/common_impl.rs:100
↓ 2 callersFunctionindex_map
(list: &Vec<U>, mapper: F)
src/common/src/collections/lang.rs:16
↓ 2 callersMethodkeyBy
(callbackFn: (value: T) => U)
typescript-api/src/stream/dataflow.ts:27
↓ 2 callersMethodkeyExtractor
(value: (val: V) => any)
typescript-api/src/connectors/connectors.ts:305
↓ 2 callersFunctionload_builder
()
src/lightflus-core/src/taskmanager/rpc.rs:34
↓ 2 callersFunctionno_found_worker
()
src/lightflus-core/src/errors/mod.rs:30
↓ 2 callersMethodprocess
(&self, message: KafkaMessage)
src/stream/src/connector.rs:321
↓ 2 callersFunctionprocessor
(message: KafkaMessage)
src/stream/tests/test_sinks.rs:113
↓ 2 callersFunctionreplace_by_env
(value: &str)
src/common/src/utils.rs:161
↓ 2 callersMethodset_multiple
( &mut self, items: &[(&K, &V)], )
src/common/src/redis.rs:47
↓ 2 callersFunctionsetup
()
src/common/tests/test_kafka.rs:10
↓ 2 callersFunctionsetup
()
src/stream/tests/test_sinks.rs:30
↓ 2 callersMethodsink
(sink: Sink<T>)
typescript-api/src/stream/dataflow.ts:37
↓ 2 callersMethodsink_event_set_to_external_and_local
( &mut self, event_set: KeyedEventSet, cx: &mut Context<'_>, )
src/stream/src/task.rs:343
↓ 2 callersMethodsplit_into_subdataflow
(&self, dataflow: &Dataflow)
src/common/src/net/cluster.rs:130
↓ 2 callersMethodstart
(&mut self, executor: StreamExecutor)
src/stream/src/task.rs:105
↓ 2 callersMethodto_kafka_message
(&self)
src/common/src/event.rs:136
↓ 2 callersMethodtopic
(value: string)
typescript-api/src/connectors/connectors.ts:47
↓ 2 callersMethodtry_write
(&self, val: T)
src/stream/src/edge.rs:54
↓ 2 callersMethodupdate_execution_id
(&mut self, execution_id: SubDataflowId)
src/common/src/net/mod.rs:160
↓ 2 callersMethodupdate_heartbeat_status
(&self, heartbeat: &Heartbeat)
src/lightflus-core/src/coordinator/executions.rs:221
↓ 2 callersMethodvalueExtractor
(value: (val: V) => any)
typescript-api/src/connectors/connectors.ts:300
↓ 1 callersMethodDescriptor
()
tools/alpha/common/probe/probe.pb.go:55
↓ 1 callersMethodDescriptor
()
tools/alpha/common/probe/probe.pb.go:101
↓ 1 callersMethodack
(&self, ack: &Ack)
src/lightflus-core/src/coordinator/scheduler.rs:52
↓ 1 callersMethodack_from_execution
(&self, ack: &Ack)
src/lightflus-core/src/coordinator/managers.rs:94
↓ 1 callersMethodack_from_task_manager
(&self, ack: Ack)
src/lightflus-core/src/coordinator/managers.rs:201
↓ 1 callersMethodadd_external_sink
(&mut self, sink: SinkImpl)
src/stream/src/task.rs:244
↓ 1 callersFunctionall_match
(elems: &Vec<T>, predicate: F)
src/common/src/collections/lang.rs:93
↓ 1 callersMethodasISink
()
typescript-api/src/connectors/connectors.ts:183
↓ 1 callersMethodasISource
()
typescript-api/src/connectors/connectors.ts:67
↓ 1 callersMethodasIStatement
()
typescript-api/src/connectors/connectors.ts:214
↓ 1 callersMethodbatch_send_event_to_operator
( &self, event_set: KeyedEventSet, )
src/lightflus-core/src/taskmanager/taskworker.rs:152
↓ 1 callersMethodbatch_write
( &self, _job_id: &Option<ResourceId>, to_operator_id: ExecutorId, _from_opera
src/stream/src/edge.rs:85
↓ 1 callersMethodbuild
(&self)
src/lightflus-core/src/taskmanager/taskworker.rs:41
↓ 1 callersMethodbuild_in_edge
Unlike out-edge which the data stream can be broadcast to multiple downstreams, each operator does have only on in-edge to receive data stream. For di
src/stream/src/task.rs:212
↓ 1 callersMethodclose_source
(&mut self)
src/stream/src/connector.rs:134
↓ 1 callersMethodcmp
(&self, other: &Self)
src/proto/src/common_impl.rs:480
↓ 1 callersMethodcmp
(&self, other: &Self)
src/common/src/types.rs:82
↓ 1 callersFunctioncreate_dataflow
( req: CreateResourceRequest, )
src/lightflus-core/src/apiserver/handler/services.rs:18
↓ 1 callersMethodcreate_dataflow
/ Attempt to deploy a new dataflow and create a JobManager. / Unless bump into network problems, JobManager will be informed the status of the deploye
src/proto/src/coordinator.rs:79
↓ 1 callersMethodcreate_sub_dataflow
/ Attempt to create a sub-dataflow
src/proto/src/taskmanager.rs:170
↓ 1 callersMethoddatabase
(db: string)
typescript-api/src/connectors/connectors.ts:259
↓ 1 callersMethoddeploy
(mut self)
src/lightflus-core/src/coordinator/executions.rs:138
↓ 1 callersMethoddeploy_dataflow
Once a dataflow is deployed, JobManager will receive the event of state transition of each subdataflow from TaskManager.
src/lightflus-core/src/coordinator/managers.rs:49
↓ 1 callersMethodevent_id
(&self)
src/common/src/event.rs:179
↓ 1 callersMethodexecute
( &'a mut self, plan: SubdataflowDeploymentPlan<'a>, )
src/lightflus-core/src/coordinator/scheduler.rs:22
↓ 1 callersMethodexecute_script
(&mut self, script: Local<'s, v8::Script>)
src/stream/src/v8_runtime.rs:97
↓ 1 callersFunctionextract_arguments_scope
( extractors: &[String], event: &LocalEvent, fn_name: &str, scope: &mut v8::HandleScope<'_, ()
src/stream/src/connector.rs:641
↓ 1 callersFunctionfile_common_probe_proto_init
()
tools/alpha/common/probe/probe.pb.go:407
↓ 1 callersFunctionfrom_pb_slice
(data: &[u8])
src/common/src/utils.rs:244
↓ 1 callersMethodgenerate_new_event_id
(&self)
src/stream/src/connector.rs:348
↓ 1 callersMethodgetCreateDataflowOptions
()
typescript-api/src/stream/context.ts:48
↓ 1 callersMethodgetDataType
()
typescript-api/src/connectors/connectors.ts:81
↓ 1 callersFunctionget_dataflow
(args: &GetResourceArgs)
src/lightflus-core/src/apiserver/handler/services.rs:47
↓ 1 callersMethodget_dataflow
/ Get the details of a dataflow. / The details contains: each operator's status, metrics, basic information, checkpoint status, etc.
src/proto/src/coordinator.rs:122
↓ 1 callersMethodget_event_time
(&self)
src/proto/src/common_impl.rs:392
↓ 1 callersMethodget_execution_id_ref
(&self)
src/proto/src/common_impl.rs:321
↓ 1 callersMethodget_host_addr
(&self)
src/proto/src/common_impl.rs:120
↓ 1 callersMethodget_host_addr_ref
(&self)
src/proto/src/common_impl.rs:127
↓ 1 callersMethodget_job_id_opt_ref
(&self)
src/proto/src/common_impl.rs:382
↓ 1 callersMethodget_kafka_partition
(&self)
src/proto/src/common_impl.rs:202
↓ 1 callersMethodget_mysql_statement
(&self)
src/proto/src/common_impl.rs:211
↓ 1 callersMethodget_nodes
(&self)
src/common/src/net/cluster.rs:246
↓ 1 callersFunctionget_opeartor_udf
(op_info: &OperatorInfo)
src/stream/src/dataflow.rs:416
↓ 1 callersMethodget_state
(&self)
src/lightflus-core/src/taskmanager/taskworker.rs:166
↓ 1 callersMethodget_states
(&self)
src/lightflus-core/src/coordinator/executions.rs:241
↓ 1 callersMethodget_sub_dataflow
Get sub dataflow states
src/proto/src/taskmanager.rs:253
↓ 1 callersFunctiongroup
(list: &Vec<N>, mut key_extractor: F)
src/common/src/collections/lang.rs:27
↓ 1 callersFunctiongroup_deque_as_btree_map
group_deque will pop all elems from deque and group them as HashMap ``` use std::collections; use common::collections::lang; let ref mut deque = coll
src/common/src/collections/lang.rs:63
↓ 1 callersMethodhas_sink
(&self)
src/proto/src/common_impl.rs:110
↓ 1 callersFunctionindex_all_match_mut
( elems: &mut Vec<T>, mut predicate: F, )
src/common/src/collections/lang.rs:101
↓ 1 callersMethodis_available
(&self)
src/common/src/net/cluster.rs:61
↓ 1 callersMethodis_dataflow_empty
(&self)
src/proto/src/apiserver_impl.rs:20
↓ 1 callersFunctionis_remote_operator
(operator: &OperatorInfo)
src/common/src/utils.rs:180
↓ 1 callersMethodis_valid
(&self)
src/proto/src/common_impl.rs:514
↓ 1 callersFunctionlocal_ip
()
src/common/src/net/mod.rs:46
↓ 1 callersFunctionmap_self
(list: &Vec<N>, mut key_extractor: F)
src/common/src/collections/lang.rs:3
↓ 1 callersMethodmsg
(&self)
src/common/src/err.rs:75
↓ 1 callersFunctionnew_key_value_state_mgt
(resource_id: &ResourceId)
src/stream/src/state.rs:15
← previousnext →101–200 of 588, ranked by callers