Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/JasonThon/lightflus
/ functions
Functions
588 in github.com/JasonThon/lightflus
⨍
Functions
588
◇
Types & classes
381
↓ 2 callers
Function
eval
( scope: &mut v8::HandleScope<'s>, code: &str, )
src/stream/src/v8_runtime.rs:498
↓ 2 callers
Method
execute
()
typescript-api/src/stream/dataflow.ts:42
↓ 2 callers
Function
execution_id_unprovided
()
src/lightflus-core/src/errors/mod.rs:19
↓ 2 callers
Function
extract_arguments
( extractors: &[String], event: &LocalEvent, fn_name: &str, )
src/stream/src/connector.rs:662
↓ 2 callers
Method
flatMap
(callbackFn: (value: T) => U[])
typescript-api/src/stream/dataflow.ts:22
↓ 2 callers
Method
from
(conn_opts: mysql_desc::ConnectionOpts)
src/common/src/db.rs:166
↓ 2 callers
Function
from_reader
(reader: R)
src/common/src/utils.rs:148
↓ 2 callers
Function
from_str
(value: &str)
src/common/src/utils.rs:157
↓ 2 callers
Method
getId
()
typescript-api/src/stream/context.ts:35
↓ 2 callers
Method
get_dataflow
( &self, job_id: &ResourceId, )
src/lightflus-core/src/coordinator/coord.rs:108
↓ 2 callers
Method
get_downstream_id_iter
(&self)
src/stream/src/task.rs:77
↓ 2 callers
Method
get_gateway
(&self)
src/common/src/net/cluster.rs:67
↓ 2 callers
Method
get_id
(&self)
src/common/src/net/cluster.rs:72
↓ 2 callers
Method
get_kafka_group
(&self)
src/proto/src/common_impl.rs:195
↓ 2 callers
Method
get_key
(&self)
src/proto/src/common_impl.rs:387
↓ 2 callers
Method
get_keyed_state
(&self, key: &[u8])
src/stream/src/state.rs:46
↓ 2 callers
Function
get_operator_state_key
(operator_id: NodeIdx, operator: &str, reference: &[u8])
src/stream/src/dataflow.rs:400
↓ 2 callers
Method
get_sub_dataflow
( &self, request: RpcRequest<ResourceId>, )
src/lightflus-core/src/taskmanager/rpc.rs:191
↓ 2 callers
Method
get_subdataflow_id
(&self)
src/proto/src/common_impl.rs:521
↓ 2 callers
Method
group
(value: string)
typescript-api/src/connectors/connectors.ts:41
↓ 2 callers
Method
has_source
(&self)
src/proto/src/common_impl.rs:100
↓ 2 callers
Function
index_map
(list: &Vec<U>, mapper: F)
src/common/src/collections/lang.rs:16
↓ 2 callers
Method
keyBy
(callbackFn: (value: T) => U)
typescript-api/src/stream/dataflow.ts:27
↓ 2 callers
Method
keyExtractor
(value: (val: V) => any)
typescript-api/src/connectors/connectors.ts:305
↓ 2 callers
Function
load_builder
()
src/lightflus-core/src/taskmanager/rpc.rs:34
↓ 2 callers
Function
no_found_worker
()
src/lightflus-core/src/errors/mod.rs:30
↓ 2 callers
Method
process
(&self, message: KafkaMessage)
src/stream/src/connector.rs:321
↓ 2 callers
Function
processor
(message: KafkaMessage)
src/stream/tests/test_sinks.rs:113
↓ 2 callers
Function
replace_by_env
(value: &str)
src/common/src/utils.rs:161
↓ 2 callers
Method
set_multiple
( &mut self, items: &[(&K, &V)], )
src/common/src/redis.rs:47
↓ 2 callers
Function
setup
()
src/common/tests/test_kafka.rs:10
↓ 2 callers
Function
setup
()
src/stream/tests/test_sinks.rs:30
↓ 2 callers
Method
sink
(sink: Sink<T>)
typescript-api/src/stream/dataflow.ts:37
↓ 2 callers
Method
sink_event_set_to_external_and_local
( &mut self, event_set: KeyedEventSet, cx: &mut Context<'_>, )
src/stream/src/task.rs:343
↓ 2 callers
Method
split_into_subdataflow
(&self, dataflow: &Dataflow)
src/common/src/net/cluster.rs:130
↓ 2 callers
Method
start
(&mut self, executor: StreamExecutor)
src/stream/src/task.rs:105
↓ 2 callers
Method
to_kafka_message
(&self)
src/common/src/event.rs:136
↓ 2 callers
Method
topic
(value: string)
typescript-api/src/connectors/connectors.ts:47
↓ 2 callers
Method
try_write
(&self, val: T)
src/stream/src/edge.rs:54
↓ 2 callers
Method
update_execution_id
(&mut self, execution_id: SubDataflowId)
src/common/src/net/mod.rs:160
↓ 2 callers
Method
update_heartbeat_status
(&self, heartbeat: &Heartbeat)
src/lightflus-core/src/coordinator/executions.rs:221
↓ 2 callers
Method
valueExtractor
(value: (val: V) => any)
typescript-api/src/connectors/connectors.ts:300
↓ 1 callers
Method
Descriptor
()
tools/alpha/common/probe/probe.pb.go:55
↓ 1 callers
Method
Descriptor
()
tools/alpha/common/probe/probe.pb.go:101
↓ 1 callers
Method
ack
(&self, ack: &Ack)
src/lightflus-core/src/coordinator/scheduler.rs:52
↓ 1 callers
Method
ack_from_execution
(&self, ack: &Ack)
src/lightflus-core/src/coordinator/managers.rs:94
↓ 1 callers
Method
ack_from_task_manager
(&self, ack: Ack)
src/lightflus-core/src/coordinator/managers.rs:201
↓ 1 callers
Method
add_external_sink
(&mut self, sink: SinkImpl)
src/stream/src/task.rs:244
↓ 1 callers
Function
all_match
(elems: &Vec<T>, predicate: F)
src/common/src/collections/lang.rs:93
↓ 1 callers
Method
asISink
()
typescript-api/src/connectors/connectors.ts:183
↓ 1 callers
Method
asISource
()
typescript-api/src/connectors/connectors.ts:67
↓ 1 callers
Method
asIStatement
()
typescript-api/src/connectors/connectors.ts:214
↓ 1 callers
Method
batch_send_event_to_operator
( &self, event_set: KeyedEventSet, )
src/lightflus-core/src/taskmanager/taskworker.rs:152
↓ 1 callers
Method
batch_write
( &self, _job_id: &Option<ResourceId>, to_operator_id: ExecutorId, _from_opera
src/stream/src/edge.rs:85
↓ 1 callers
Method
build
(&self)
src/lightflus-core/src/taskmanager/taskworker.rs:41
↓ 1 callers
Method
build_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 callers
Method
close_source
(&mut self)
src/stream/src/connector.rs:134
↓ 1 callers
Method
cmp
(&self, other: &Self)
src/proto/src/common_impl.rs:480
↓ 1 callers
Method
cmp
(&self, other: &Self)
src/common/src/types.rs:82
↓ 1 callers
Function
create_dataflow
( req: CreateResourceRequest, )
src/lightflus-core/src/apiserver/handler/services.rs:18
↓ 1 callers
Method
create_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 callers
Method
create_sub_dataflow
/ Attempt to create a sub-dataflow
src/proto/src/taskmanager.rs:170
↓ 1 callers
Method
database
(db: string)
typescript-api/src/connectors/connectors.ts:259
↓ 1 callers
Method
deploy
(mut self)
src/lightflus-core/src/coordinator/executions.rs:138
↓ 1 callers
Method
deploy_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 callers
Method
event_id
(&self)
src/common/src/event.rs:179
↓ 1 callers
Method
execute
( &'a mut self, plan: SubdataflowDeploymentPlan<'a>, )
src/lightflus-core/src/coordinator/scheduler.rs:22
↓ 1 callers
Method
execute_script
(&mut self, script: Local<'s, v8::Script>)
src/stream/src/v8_runtime.rs:97
↓ 1 callers
Function
extract_arguments_scope
( extractors: &[String], event: &LocalEvent, fn_name: &str, scope: &mut v8::HandleScope<'_, ()
src/stream/src/connector.rs:641
↓ 1 callers
Function
file_common_probe_proto_init
()
tools/alpha/common/probe/probe.pb.go:407
↓ 1 callers
Function
from_pb_slice
(data: &[u8])
src/common/src/utils.rs:244
↓ 1 callers
Method
generate_new_event_id
(&self)
src/stream/src/connector.rs:348
↓ 1 callers
Method
getCreateDataflowOptions
()
typescript-api/src/stream/context.ts:48
↓ 1 callers
Method
getDataType
()
typescript-api/src/connectors/connectors.ts:81
↓ 1 callers
Function
get_dataflow
(args: &GetResourceArgs)
src/lightflus-core/src/apiserver/handler/services.rs:47
↓ 1 callers
Method
get_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 callers
Method
get_event_time
(&self)
src/proto/src/common_impl.rs:392
↓ 1 callers
Method
get_execution_id_ref
(&self)
src/proto/src/common_impl.rs:321
↓ 1 callers
Method
get_host_addr
(&self)
src/proto/src/common_impl.rs:120
↓ 1 callers
Method
get_host_addr_ref
(&self)
src/proto/src/common_impl.rs:127
↓ 1 callers
Method
get_job_id_opt_ref
(&self)
src/proto/src/common_impl.rs:382
↓ 1 callers
Method
get_kafka_partition
(&self)
src/proto/src/common_impl.rs:202
↓ 1 callers
Method
get_mysql_statement
(&self)
src/proto/src/common_impl.rs:211
↓ 1 callers
Method
get_nodes
(&self)
src/common/src/net/cluster.rs:246
↓ 1 callers
Function
get_opeartor_udf
(op_info: &OperatorInfo)
src/stream/src/dataflow.rs:416
↓ 1 callers
Method
get_state
(&self)
src/lightflus-core/src/taskmanager/taskworker.rs:166
↓ 1 callers
Method
get_states
(&self)
src/lightflus-core/src/coordinator/executions.rs:241
↓ 1 callers
Method
get_sub_dataflow
Get sub dataflow states
src/proto/src/taskmanager.rs:253
↓ 1 callers
Function
group
(list: &Vec<N>, mut key_extractor: F)
src/common/src/collections/lang.rs:27
↓ 1 callers
Function
group_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 callers
Method
has_sink
(&self)
src/proto/src/common_impl.rs:110
↓ 1 callers
Function
index_all_match_mut
( elems: &mut Vec<T>, mut predicate: F, )
src/common/src/collections/lang.rs:101
↓ 1 callers
Method
is_available
(&self)
src/common/src/net/cluster.rs:61
↓ 1 callers
Method
is_dataflow_empty
(&self)
src/proto/src/apiserver_impl.rs:20
↓ 1 callers
Function
is_remote_operator
(operator: &OperatorInfo)
src/common/src/utils.rs:180
↓ 1 callers
Method
is_valid
(&self)
src/proto/src/common_impl.rs:514
↓ 1 callers
Function
local_ip
()
src/common/src/net/mod.rs:46
↓ 1 callers
Function
map_self
(list: &Vec<N>, mut key_extractor: F)
src/common/src/collections/lang.rs:3
↓ 1 callers
Method
msg
(&self)
src/common/src/err.rs:75
↓ 1 callers
Function
new_key_value_state_mgt
(resource_id: &ResourceId)
src/stream/src/state.rs:15
← previous
next →
101–200 of 588, ranked by callers