Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/apache/storm
/ functions
Functions
27,770 in github.com/apache/storm
⨍
Functions
27,770
◇
Types & classes
4,363
↓ 1 callers
Method
addShutdownHooks
(MyThd... threads)
examples/storm-perf/src/main/java/org/apache/storm/perf/toolstest/JCToolsPerfTest.java:95
↓ 1 callers
Method
addSink
(TopologyBuilder topologyBuilder, SinkNode sinkNode)
storm-client/src/jvm/org/apache/storm/streams/StreamBuilder.java:447
↓ 1 callers
Method
addSourcedStateNode
(List<Stream> sources, Node newNode)
storm-client/src/jvm/org/apache/storm/trident/TridentTopology.java:1006
↓ 1 callers
Method
addSpout
(TopologyBuilder topologyBuilder, SpoutNode spout)
storm-client/src/jvm/org/apache/storm/streams/StreamBuilder.java:443
↓ 1 callers
Method
addSpout
(String id, IRichSpout spout)
flux/flux-core/src/main/java/org/apache/storm/flux/model/ExecutionContext.java:58
↓ 1 callers
Method
addSpoutAggStats
If aggStats are not populated, compute common and component(spout) agg and create placeholder stat. This allow the topology page to show component spe
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:4508
↓ 1 callers
Method
addStatefulBolt
(TopologyBuilder topologyBuilder, String boltId,
storm-client/src/jvm/org/apache/storm/streams/StreamBuilder.java:475
↓ 1 callers
Method
addStreamToInitialProcessors
(Multimap<String, ProcessorNode> streamToInitialProcessors)
storm-client/src/jvm/org/apache/storm/streams/StatefulProcessorBolt.java:83
↓ 1 callers
Method
addSystemComponents
(Map<String, Object> conf, StormTopology topology)
storm-client/src/jvm/org/apache/storm/daemon/StormCommon.java:426
↓ 1 callers
Method
addSystemStreams
(StormTopology topology)
storm-client/src/jvm/org/apache/storm/daemon/StormCommon.java:318
↓ 1 callers
Method
addTask
add task to cgroup. @param taskid task id of task to add
storm-client/src/jvm/org/apache/storm/container/cgroup/CgroupCommonOperation.java:25
↓ 1 callers
Method
addTaskHooks
()
storm-client/src/jvm/org/apache/storm/daemon/Task.java:294
↓ 1 callers
Method
addTimer
(String name, Timer timer)
storm-client/src/jvm/org/apache/storm/metrics2/TaskMetricRepo.java:53
↓ 1 callers
Method
addToCache
(String key, String value)
examples/storm-starter/src/jvm/org/apache/storm/starter/ResourceAwareExampleTopology.java:115
↓ 1 callers
Method
addToLeaderLockQueue
queue up for leadership lock. The call returns immediately and the caller must check isLeader() to perform any leadership action. This method can be c
storm-client/src/jvm/org/apache/storm/nimbus/ILeaderElector.java:36
↓ 1 callers
Method
addToLeaderLockQueue
()
storm-server/src/main/java/org/apache/storm/zookeeper/LeaderElectorImp.java:54
↓ 1 callers
Method
addToWindowManager
(int tupleIndex, String effectiveBatchId, TridentTuple tridentTuple)
storm-client/src/jvm/org/apache/storm/trident/windowing/StoreBasedTridentWindowManager.java:144
↓ 1 callers
Method
addTopoConf
Add a new topology config. @param topoId the id of the topology @param who who is doing it @param topoConf the topology conf itself @throws Authorizat
storm-server/src/main/java/org/apache/storm/daemon/nimbus/TopoCache.java:187
↓ 1 callers
Method
addTopoToHistoryLog
(String topoId, Map<String, Object> topoConf)
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:2944
↓ 1 callers
Method
addTopology
Add a topology to the pool. @param td the topology to add
storm-server/src/main/java/org/apache/storm/scheduler/multitenant/NodePool.java:128
↓ 1 callers
Method
addTopology
Add a new topology. @param topoId the id of the topology @param who who is doing it @param topo the topology itself @throws AuthorizationException if
storm-server/src/main/java/org/apache/storm/daemon/nimbus/TopoCache.java:104
↓ 1 callers
Method
addTopologyHistory
(LSTopoHistory lsTopoHistory)
storm-client/src/jvm/org/apache/storm/utils/LocalState.java:210
↓ 1 callers
Method
addTriggerField
(Fields functionFields)
storm-client/src/jvm/org/apache/storm/trident/Stream.java:757
↓ 1 callers
Method
addTuplesBatch
Add received batch of tuples to cache/store and add them to {@code WindowManager}.
storm-client/src/jvm/org/apache/storm/trident/windowing/ITridentWindowManager.java:40
↓ 1 callers
Method
addV2Metrics
(int taskId, List<IMetricsConsumer.DataPoint> dataPoints, int interval)
storm-client/src/jvm/org/apache/storm/executor/Executor.java:368
↓ 1 callers
Method
addWindowedBolt
(TopologyBuilder topologyBuilder, String boltId,
storm-client/src/jvm/org/apache/storm/streams/StreamBuilder.java:508
↓ 1 callers
Function
add_tag_to_dicts
(hash_to_tags, tag_to_hash, tag, manifest_hash, comment)
bin/docker-to-squash.py:470
↓ 1 callers
Method
add_to_executors
(ExecutorInfo elem)
storm-client/src/jvm/org/apache/storm/generated/SupervisorWorkerHeartbeat.java:206
↓ 1 callers
Method
add_to_metrics
(WorkerMetricPoint elem)
storm-client/src/jvm/org/apache/storm/generated/WorkerMetricList.java:153
↓ 1 callers
Method
add_to_supervisor_summaries
(SupervisorSummary elem)
storm-client/src/jvm/org/apache/storm/generated/SupervisorPageInfo.java:163
↓ 1 callers
Method
add_to_worker_hooks
(java.nio.ByteBuffer elem)
storm-client/src/jvm/org/apache/storm/generated/StormTopology.java:434
↓ 1 callers
Method
add_to_worker_summaries
(WorkerSummary elem)
storm-client/src/jvm/org/apache/storm/generated/SupervisorPageInfo.java:203
↓ 1 callers
Method
adjustResourcesForEvictedTopology
(Cluster cluster, TopologyDetails evict)
storm-server/src/main/java/org/apache/storm/scheduler/resource/ResourceAwareScheduler.java:376
↓ 1 callers
Method
advanceHead
()
examples/storm-starter/src/jvm/org/apache/storm/starter/tools/SlidingWindowCounter.java:107
↓ 1 callers
Method
aggBoltExecWinStats
aggregate windowed stats from a bolt executor stats with a Map of accumulated stats.
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:946
↓ 1 callers
Method
aggBoltStreamsLatAndCount
aggregate number executed and process & execute latencies.
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:149
↓ 1 callers
Method
aggCompExecStats
Combines the aggregate stats of one executor with the given map, selecting the appropriate window and including system components as specified.
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:1096
↓ 1 callers
Method
aggCompExecsStats
aggregate component executor stats. @param exec2hostPort a Map of {executor -> host+port} @param task2component a Map of {task id -> component} @par
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:1205
↓ 1 callers
Method
aggPartition
(S stream)
storm-client/src/jvm/org/apache/storm/trident/fluent/GlobalAggregationScheme.java:19
↓ 1 callers
Method
aggPreMergeCompPageBolt
pre-merge component page bolt stats from an executor heartbeat 1. computes component capacity 2. converts map keys of stats 3. filters streams if nece
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:195
↓ 1 callers
Method
aggPreMergeCompPageSpout
pre-merge component page spout stats from an executor heartbeat 1. computes component capacity 2. converts map keys of stats 3. filters streams if nec
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:254
↓ 1 callers
Method
aggPreMergeTopoPageBolt
pre-merge component stats of specified bolt id. @param beat executor heartbeat data @param window specified window @param includeSys whethe
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:304
↓ 1 callers
Method
aggPreMergeTopoPageSpout
pre-merge component stats of specified spout id and returns { comp id -> comp-stats }.
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:345
↓ 1 callers
Method
aggSpoutExecWinStats
aggregate windowed stats from a spout executor stats with a Map of accumulated stats.
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:990
↓ 1 callers
Method
aggSpoutStreamsLatAndCount
Aggregates number acked and complete latencies.
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:171
↓ 1 callers
Method
aggTopoExecStats
A helper function that does the common work to aggregate stats of one executor with the given map for the topology page.
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:523
↓ 1 callers
Method
aggregate
(Aggregator agg, Fields functionFields)
storm-client/src/jvm/org/apache/storm/trident/fluent/GroupedStream.java:49
↓ 1 callers
Method
aggregate
Aggregates the values in this stream using the aggregator. This does a global aggregation of values across all partitions. <p> If the stream is window
storm-client/src/jvm/org/apache/storm/streams/Stream.java:195
↓ 1 callers
Method
aggregateCompStats
Aggregate the stats for a component over a given window of time.
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:1059
↓ 1 callers
Method
aggregateSpoutStats
aggregate spout stats. @param statsSeq a seq of ExecutorStats @param includeSys whether to include system streams @return aggregated spout stats: {
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:746
↓ 1 callers
Method
aggregateSpoutStreams
aggregate all spout streams. @param stats a Map of {metric -> win -> stream id -> value} @return a Map of {metric -> win -> aggregated value}
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:911
↓ 1 callers
Method
aggregateTopoStats
(String win, boolean includeSys, List<Map<String, Object>> heartbeats)
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:622
↓ 1 callers
Method
aliveExecutors
(String topoId, Set<List<Integer>> allExecutors, Assignment assignment)
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:2103
↓ 1 callers
Method
all
()
storm-client/src/jvm/org/apache/storm/trident/topology/TridentBoltExecutor.java:310
↓ 1 callers
Method
all
()
storm-client/src/jvm/org/apache/storm/streams/GroupingInfo.java:57
↓ 1 callers
Method
allTasksReady
()
storm-server/src/main/java/org/apache/storm/nimbus/NimbusHeartbeatsPressureTest.java:106
↓ 1 callers
Method
allocatedTopologies
(Map<String, Set<Set<ExecutorDetails>>> topologyToWorkerSpecs)
storm-server/src/main/java/org/apache/storm/scheduler/IsolationScheduler.java:353
↓ 1 callers
Method
anonymize
Create a new version of this topology with identifiable information removed. @return the anonymized version of the TopologyLoadConf.
examples/storm-loadgen/src/main/java/org/apache/storm/loadgen/TopologyLoadConf.java:331
↓ 1 callers
Method
anonymizeTopoConf
(Map<String, Object> topoConf)
examples/storm-loadgen/src/main/java/org/apache/storm/loadgen/TopologyLoadConf.java:387
↓ 1 callers
Method
anyNonCpuOverZero
Are any of the non cpu resources positive. @return true of any of the non cpu resources are positive. False if they are all <= 0.
storm-server/src/main/java/org/apache/storm/scheduler/resource/normalization/NormalizedResources.java:443
↓ 1 callers
Method
apply
Applies the `Assembly` to a given {@link org.apache.storm.trident.Stream}.
storm-client/src/jvm/org/apache/storm/trident/operation/Assembly.java:32
↓ 1 callers
Method
apply
(ConsumerRecord<K, V> record)
external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/DefaultRecordTranslator.java:30
↓ 1 callers
Method
applyOn
(TopologyContext topologyContext)
storm-client/src/jvm/org/apache/storm/hooks/info/BoltFailInfo.java:32
↓ 1 callers
Method
applyOn
(TopologyContext topologyContext)
storm-client/src/jvm/org/apache/storm/hooks/info/BoltAckInfo.java:32
↓ 1 callers
Method
applyProperties
(ObjectDef bean, Object instance, ExecutionContext context)
flux/flux-core/src/main/java/org/apache/storm/flux/FluxBuilder.java:277
↓ 1 callers
Method
applyProxy
()
storm-submit-tools/src/main/java/org/apache/storm/submit/dependency/DependencyResolver.java:132
↓ 1 callers
Function
applyThemeToNetwork
()
storm-webapp/src/main/webapp/js/visualization.js:159
↓ 1 callers
Method
applyUUIDToFileName
(String fileName)
storm-client/src/jvm/org/apache/storm/dependency/DependencyBlobStoreUtils.java:33
↓ 1 callers
Method
areAllConnectionsReady
()
storm-client/src/jvm/org/apache/storm/daemon/worker/WorkerState.java:675
↓ 1 callers
Method
areAllSupervisorsWaiting
()
storm-server/src/main/java/org/apache/storm/LocalCluster.java:787
↓ 1 callers
Method
areAllWorkersWaiting
()
storm-server/src/main/java/org/apache/storm/LocalCluster.java:303
↓ 1 callers
Method
areAnyOverZero
(boolean skipCpuCheck)
storm-server/src/main/java/org/apache/storm/scheduler/resource/normalization/NormalizedResources.java:418
↓ 1 callers
Method
areWorkerTokensSupported
Check if worker tokens are supported by this transport. @return true if they are else false.
storm-client/src/jvm/org/apache/storm/security/auth/ITransportPlugin.java:65
↓ 1 callers
Method
asAccessControl
(String param)
storm-core/src/jvm/org/apache/storm/command/Blobstore.java:319
↓ 1 callers
Method
asPath
(String... parts)
storm-server/src/test/java/org/apache/storm/daemon/supervisor/ContainerTest.java:62
↓ 1 callers
Method
assertOneElement
(Collection<?> collection)
integration-test/src/test/java/org/apache/storm/st/utils/AssertUtil.java:53
↓ 1 callers
Method
assertTopologiesBeenEvicted
(Cluster cluster, Class strategyClass, Set<String> evictedTopologies, String... topoNames)
storm-server/src/test/java/org/apache/storm/scheduler/resource/TestUtilsForResourceAwareScheduler.java:523
↓ 1 callers
Method
assertTwoElements
(Collection<?> collection)
integration-test/src/test/java/org/apache/storm/st/utils/AssertUtil.java:62
↓ 1 callers
Method
assignSingleBoundAcker
<p> Remove the head of unassigned ackers and attempt to assign it to a workerSlot as a bound acker. </p> @param node RasNode on which to sch
storm-server/src/main/java/org/apache/storm/scheduler/resource/strategies/scheduling/SchedulingSearcherState.java:352
↓ 1 callers
Method
assigned
(Collection<Integer> ports)
storm-server/src/main/java/org/apache/storm/scheduler/ISupervisor.java:45
↓ 1 callers
Method
assignmentsForHost
Pick out assignments for a specific host from all assignments. This could include multiple NUMA supervisors on an individual host. @param assignmentM
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:1690
↓ 1 callers
Method
assignmentsForNodeId
Pick out assignments for specific NodeId from all assignments. @param assignmentMap stormId -> assignment map @param nodeId supervisor node id
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:1709
↓ 1 callers
Method
auditAssignmentChanges
Check new assignments with existing assignments and determine difference is any. @param existingAssignments non-null map of topology-id to existing a
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:843
↓ 1 callers
Method
authenticated
(Channel c)
storm-client/src/jvm/org/apache/storm/messaging/netty/Server.java:292
↓ 1 callers
Method
awaitTermination
()
storm-webapp/src/main/java/org/apache/storm/daemon/logviewer/LogviewerServer.java:137
↓ 1 callers
Method
backPressureWaitStrategy
()
storm-client/src/jvm/org/apache/storm/executor/spout/SpoutExecutor.java:226
↓ 1 callers
Method
backtrack
Backtrack to prior executor that was directly assigned. This excludes bound-ackers. @param execToComp map from executor to component. @param nodesFor
storm-server/src/main/java/org/apache/storm/scheduler/resource/strategies/scheduling/SchedulingSearcherState.java:284
↓ 1 callers
Method
badSlots
(Map<WorkerSlot, List<ExecutorDetails>> existingSlots, int numExecutors, int numWorkers)
storm-server/src/main/java/org/apache/storm/scheduler/DefaultScheduler.java:33
↓ 1 callers
Method
badSlots
(SupervisorDetails supervisor, String supervisorKey)
storm-server/src/main/java/org/apache/storm/scheduler/blacklist/BlacklistScheduler.java:177
↓ 1 callers
Method
bitXorVals
(List<Long> coll)
storm-client/src/jvm/org/apache/storm/utils/Utils.java:318
↓ 1 callers
Method
blobChanging
Informs the listener that a blob has changed and is ready to update and replace a localized blob that has been marked as tied to the life cycle of the
storm-server/src/main/java/org/apache/storm/localizer/BlobChangingCallback.java:39
↓ 1 callers
Method
blobExists
Checks if a blob exists. @param key blobstore key @param who subject @throws AuthorizationException if authorization is failed
external/storm-hdfs-blobstore/src/main/java/org/apache/storm/hdfs/blobstore/HdfsBlobStore.java:304
↓ 1 callers
Method
blobNeedsWorkerRestart
Given the blob information returns the value of the workerRestart field, handling it being a boolean value, or if it's not specified then returns fals
storm-server/src/main/java/org/apache/storm/daemon/supervisor/SupervisorUtils.java:137
↓ 1 callers
Method
blobstoreMapToLocalresources
Returns a list of LocalResources based on the blobstore-map passed in.
storm-server/src/main/java/org/apache/storm/daemon/supervisor/SupervisorUtils.java:144
↓ 1 callers
Method
blobstoreMaxKeySequenceNumberPath
(String key)
storm-client/src/jvm/org/apache/storm/cluster/ClusterUtils.java:143
↓ 1 callers
Method
bolt
(BoltAggregateStats value)
storm-client/src/jvm/org/apache/storm/generated/SpecificAggregateStats.java:125
↓ 1 callers
Method
bolt
(BoltStats value)
storm-client/src/jvm/org/apache/storm/generated/ExecutorSpecificStats.java:125
↓ 1 callers
Method
boltAck
(BoltAckInfo info)
storm-client/src/jvm/org/apache/storm/hooks/ITaskHook.java:37
↓ 1 callers
Method
boltExecute
(List<Tuple> tuples, List<Tuple> newTuples, List<Tuple> expiredTuples, Long timestamp)
storm-client/src/jvm/org/apache/storm/topology/WindowedBoltExecutor.java:370
↓ 1 callers
Method
boltFail
(BoltFailInfo info)
storm-client/src/jvm/org/apache/storm/hooks/ITaskHook.java:39
← previous
next →
6,601–6,700 of 27,770, ranked by callers