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
cleanup
Called during shutdown.
storm-server/src/main/java/org/apache/storm/nimbus/ITopologyActionNotifierPlugin.java:35
↓ 1 callers
Method
cleanup
()
external/storm-jdbc/src/main/java/org/apache/storm/jdbc/common/HikariCPConnectionProvider.java:60
↓ 1 callers
Method
cleanupChildOpts
(Object value)
examples/storm-loadgen/src/main/java/org/apache/storm/loadgen/TopologyLoadConf.java:404
↓ 1 callers
Method
cleanupEmptyTopoDirectory
Delete the topo dir if it contains zero port dirs.
storm-webapp/src/main/java/org/apache/storm/daemon/logviewer/utils/LogCleaner.java:236
↓ 1 callers
Method
cleanupOrphanedData
Clean up any temporary files. This will be called after updating a blob, either successfully or if an error has occurred. The goal is to find any fil
storm-server/src/main/java/org/apache/storm/localizer/LocallyCachedBlob.java:166
↓ 1 callers
Method
cleanupStats
()
storm-client/src/jvm/org/apache/storm/stats/CommonStats.java:64
↓ 1 callers
Function
cleanup_handle_dead_images
(dead_images, known_images)
bin/docker-to-squash.py:1506
↓ 1 callers
Function
cleanup_handle_stale_images
(stale_images)
bin/docker-to-squash.py:1492
↓ 1 callers
Function
cleanup_handle_tagged_images
(tagged_images)
bin/docker-to-squash.py:1466
↓ 1 callers
Function
cleanup_handle_untagged_images
(untagged_images)
bin/docker-to-squash.py:1479
↓ 1 callers
Method
clearHashedInputs
()
storm-client/src/jvm/org/apache/storm/bolt/JoinBolt.java:188
↓ 1 callers
Method
clearIteratorPins
()
storm-client/src/jvm/org/apache/storm/windowing/persistence/WindowState.java:150
↓ 1 callers
Method
clearModified
()
storm-client/src/jvm/org/apache/storm/windowing/persistence/WindowState.java:367
↓ 1 callers
Method
clearPunctuationState
(ProcessorNode processorNode)
storm-client/src/jvm/org/apache/storm/streams/ProcessorBoltDelegate.java:312
↓ 1 callers
Method
clearRecoveryState
(TaskStream stream)
storm-client/src/jvm/org/apache/storm/topology/StatefulWindowedBoltExecutor.java:128
↓ 1 callers
Method
clearState
(String id)
storm-client/src/jvm/org/apache/storm/utils/RegisteredGlobalState.java:51
↓ 1 callers
Method
clearTrackId
()
storm-client/src/jvm/org/apache/storm/executor/LocalExecutor.java:46
↓ 1 callers
Method
cloneKerberosTicket
(final KerberosTicket kerberosTicket)
storm-client/src/jvm/org/apache/storm/security/auth/ClientAuthUtils.java:553
↓ 1 callers
Method
close
()
storm-client/src/jvm/org/apache/storm/security/auth/ThriftClient.java:222
↓ 1 callers
Method
close
()
storm-client/src/jvm/org/apache/storm/utils/BufferFileInputStream.java:46
↓ 1 callers
Method
close
()
storm-client/src/jvm/org/apache/storm/pacemaker/PacemakerClient.java:265
↓ 1 callers
Method
close
()
storm-server/src/main/java/org/apache/storm/LocalDRPC.java:91
↓ 1 callers
Method
close
()
examples/storm-loadgen/src/main/java/org/apache/storm/loadgen/ScopedTopologySet.java:81
↓ 1 callers
Method
close
()
external/storm-hdfs-blobstore/src/main/java/org/apache/storm/hdfs/blobstore/HdfsClientBlobStore.java:130
↓ 1 callers
Method
closeChannel
()
storm-client/src/jvm/org/apache/storm/messaging/netty/Client.java:507
↓ 1 callers
Method
closeFlushScheduler
()
storm-client/src/jvm/org/apache/storm/metric/FileBasedEventLogger.java:206
↓ 1 callers
Method
closeProducer
()
external/storm-kafka-client/src/test/java/org/apache/storm/kafka/KafkaUnit.java:120
↓ 1 callers
Method
closeResources
()
storm-client/src/jvm/org/apache/storm/daemon/worker/WorkerState.java:669
↓ 1 callers
Method
clusterCapacity
(Collection<SupervisorDistribution> supervisorDistributions)
storm-server/src/test/java/org/apache/storm/scheduler/resource/strategies/scheduling/TestLargeCluster.java:561
↓ 1 callers
Method
coGroupByKeyPartition
(PairStream<K, V1> otherStream)
storm-client/src/jvm/org/apache/storm/streams/PairStream.java:396
↓ 1 callers
Method
collectExecutorStats
(Set<MetricsItem> items)
examples/storm-perf/src/main/java/org/apache/storm/perf/utils/BasicMetricsCollector.java:218
↓ 1 callers
Method
collectMetricsAndKill
(String topologyName, Integer pollInterval, int duration)
examples/storm-perf/src/main/java/org/apache/storm/perf/utils/Helper.java:47
↓ 1 callers
Method
collectSpoutLatency
(Set<MetricsItem> items)
examples/storm-perf/src/main/java/org/apache/storm/perf/utils/BasicMetricsCollector.java:234
↓ 1 callers
Method
collectSpoutThroughput
(Set<MetricsItem> items)
examples/storm-perf/src/main/java/org/apache/storm/perf/utils/BasicMetricsCollector.java:229
↓ 1 callers
Method
collectThroughput
(Set<MetricsItem> items)
examples/storm-perf/src/main/java/org/apache/storm/perf/utils/BasicMetricsCollector.java:224
↓ 1 callers
Method
collectTopologyStats
(Set<MetricsItem> items)
examples/storm-perf/src/main/java/org/apache/storm/perf/utils/BasicMetricsCollector.java:213
↓ 1 callers
Method
collectTridentTupleOrKey
(TridentBatchTuple tridentBatchTuple, List<String> keys)
storm-client/src/jvm/org/apache/storm/trident/windowing/StoreBasedTridentWindowManager.java:186
↓ 1 callers
Method
combine
(T a, T b)
storm-client/src/jvm/org/apache/storm/metric/api/ICombiner.java:18
↓ 1 callers
Method
combinePartition
(CombinerAggregator<? super T, A, ?> aggregator)
storm-client/src/jvm/org/apache/storm/streams/Stream.java:442
↓ 1 callers
Method
combinePartition
(CombinerAggregator<? super V, A, ?> aggregator)
storm-client/src/jvm/org/apache/storm/streams/PairStream.java:442
↓ 1 callers
Method
commandFilePath
(String dir, String commandTag)
storm-server/src/main/java/org/apache/storm/container/oci/OciContainerManager.java:150
↓ 1 callers
Method
commit
Marks an offset as committed. This method has side effects - it sets the internal state in such a way that future calls to {@link #findNextCommitOffse
external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/internal/OffsetManager.java:170
↓ 1 callers
Method
commitNewVersion
Commit the new version and make it available for the end user. PRECONDITION: uncompressToTempLocationIfNeeded will have been called. PRECONDITION: thi
storm-server/src/main/java/org/apache/storm/localizer/LocallyCachedBlob.java:159
↓ 1 callers
Method
commitOffsetForPartitionZeroOnly
()
external/storm-kafka-monitor/src/test/java/org/apache/storm/kafka/monitor/KafkaOffsetLagUtilTest.java:103
↓ 1 callers
Method
committerBatches
(Group g, Map<Node, String> batchGroupMap)
storm-client/src/jvm/org/apache/storm/trident/TridentTopology.java:323
↓ 1 callers
Method
compactWindow
expires events that fall out of the window every EXPIRE_EVENTS_THRESHOLD so that the window does not grow too big.
storm-client/src/jvm/org/apache/storm/windowing/WindowManager.java:174
↓ 1 callers
Method
compare
(T value1, T value2)
storm-client/src/jvm/org/apache/storm/trident/operation/builtin/ComparisonAggregator.java:33
↓ 1 callers
Method
compare
(TridentTuple tuple1, TridentTuple tuple2)
examples/storm-starter/src/jvm/org/apache/storm/starter/trident/TridentMinMaxOfDevicesTopology.java:115
↓ 1 callers
Method
compare
(TridentTuple tuple1, TridentTuple tuple2)
examples/storm-starter/src/jvm/org/apache/storm/starter/trident/TridentMinMaxOfVehiclesTopology.java:95
↓ 1 callers
Method
compareTo
(Bolt other)
storm-client/src/jvm/org/apache/storm/generated/Bolt.java:301
↓ 1 callers
Method
compareTo
(MessageId rhs)
external/storm-hdfs/src/main/java/org/apache/storm/hdfs/spout/HdfsSpout.java:807
↓ 1 callers
Method
complete
()
storm-server/src/main/java/org/apache/storm/localizer/TimePortAndAssignment.java:47
↓ 1 callers
Method
completeDrpc
(DefaultDirectedGraph<Node, IndexedEdge> graph, Map<String, List<Node>> colocate, UniqueIdGen gen)
storm-client/src/jvm/org/apache/storm/trident/TridentTopology.java:179
↓ 1 callers
Method
completed
()
storm-client/src/jvm/org/apache/storm/testing/TestEventLogSpout.java:87
↓ 1 callers
Method
completedWaitingSubprocess
()
storm-client/src/jvm/org/apache/storm/spout/ShellSpout.java:265
↓ 1 callers
Method
completelyRemove
Completely remove anything that is cached locally for this blob and all tracking files also stored for it. This will be called after the blob was dete
storm-server/src/main/java/org/apache/storm/localizer/LocallyCachedBlob.java:173
↓ 1 callers
Method
completelyRemoveUnusedUser
(Path localBaseDir, String user)
storm-server/src/main/java/org/apache/storm/localizer/LocalizedResource.java:149
↓ 1 callers
Method
computeAssignedGenericResources
(Collection<WorkerResources> workers)
storm-server/src/main/java/org/apache/storm/daemon/nimbus/TopologyResources.java:92
↓ 1 callers
Method
computeBoltCapacity
computes max bolt capacity. @param executorSumms a list of ExecutorSummary @return max bolt capacity
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:1574
↓ 1 callers
Method
computeClusterCpuResource
()
storm-server/src/main/java/org/apache/storm/scheduler/Cluster.java:939
↓ 1 callers
Method
computeClusterGenericResources
()
storm-server/src/main/java/org/apache/storm/scheduler/Cluster.java:963
↓ 1 callers
Method
computeClusterMemoryResource
()
storm-server/src/main/java/org/apache/storm/scheduler/Cluster.java:952
↓ 1 callers
Method
computeComponentConstraints
()
storm-server/src/main/java/org/apache/storm/scheduler/resource/strategies/scheduling/ConstraintSolverConfig.java:71
↓ 1 callers
Method
computeComponentMap
Compute a map of user topology specific components. @param topology userTopology @return a map of Component Id to Component
storm-server/src/main/java/org/apache/storm/scheduler/TopologyDetails.java:225
↓ 1 callers
Method
computeExecutorCapacity
Compute the capacity of a executor. approximation of the % of time spent doing real work. @param summary the stats for the executor. @return the capac
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:1590
↓ 1 callers
Method
computeExecutorToComponent
(String topoId, StormBase base, Map<String,
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:2136
↓ 1 callers
Method
computeMD5
(List<File> files)
storm-buildtools/storm-maven-plugins/src/main/java/org/apache/storm/maven/plugin/versioninfo/VersionInfoMojo.java:273
↓ 1 callers
Method
computeMaxSchedulingTimeMs
(Map<String, Object> topoConf)
storm-server/src/main/java/org/apache/storm/scheduler/resource/strategies/scheduling/BaseResourceAwareStrategy.java:237
↓ 1 callers
Method
computeNewSchedulerAssignments
(Map<String, Assignment> existingAssignments,
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:2335
↓ 1 callers
Method
computeSupervisorToDeadPorts
(Map<String, Assignment> existingAssignments,
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:2193
↓ 1 callers
Method
computeTopoToExecToNodePort
convert {topology-id -> SchedulerAssignment} to {topology-id -> {executor [node port]}}. @return {topology-id -> {executor [node port]}} mapping
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:795
↓ 1 callers
Method
computeTopoToNodePortToResources
Convert {topology-id -> SchedulerAssignment} to {topology-id -> {WorkerSlot WorkerResources}}. Make sure this can deal with other non-RAS schedulers l
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:827
↓ 1 callers
Method
computeTopologyToAliveExecutors
compute a topology-id -> alive executors map. @param existingAssignment the current assignments @param topologyToExecutors the executors for the cur
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:2174
↓ 1 callers
Method
computeTopologyToExecutors
(Map<String, StormBase> bases)
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:2148
↓ 1 callers
Method
computeTopologyToSchedulerAssignment
Convert assignment information in zk to SchedulerAssignment, so it can be used by scheduler api. @param existingAssignments current assignments
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:2228
↓ 1 callers
Method
computeWaterMarkTs
Computes the min ts across all streams.
storm-client/src/jvm/org/apache/storm/windowing/WaterMarkEventGenerator.java:100
↓ 1 callers
Method
computeWorkerSpecs
(TopologyDetails topology)
storm-server/src/main/java/org/apache/storm/scheduler/IsolationScheduler.java:217
↓ 1 callers
Method
config
()
storm-server/src/main/java/org/apache/storm/scheduler/resource/ResourceAwareScheduler.java:108
↓ 1 callers
Method
config
()
storm-server/src/main/java/org/apache/storm/scheduler/multitenant/MultitenantScheduler.java:82
↓ 1 callers
Method
configureMetricProcessor
Configures metric processor (running on supervisor) to use the class specified in the conf. @param conf the supervisor config @return WorkerMetricsPr
storm-server/src/main/java/org/apache/storm/metricstore/MetricStoreConfig.java:52
↓ 1 callers
Method
confirmAssigned
(int port)
storm-server/src/main/java/org/apache/storm/scheduler/ISupervisor.java:39
↓ 1 callers
Method
connect
Connect to the specified server via framed transport. @param transport The underlying Thrift transport. @param serverHost server host @param asUser
storm-client/src/jvm/org/apache/storm/security/auth/ITransportPlugin.java:51
↓ 1 callers
Method
connectToFixedPort
(Map<String, Object> stormConf, int port)
storm-core/test/jvm/org/apache/storm/messaging/netty/NettyTest.java:378
↓ 1 callers
Method
connectionEstablished
(Channel channel)
storm-client/src/jvm/org/apache/storm/messaging/netty/Server.java:210
↓ 1 callers
Method
connectionFactory
Provides the JMS <code>ConnectionFactory</code> @return the connection factory @throws Exception
external/storm-jms/src/test/java/org/apache/storm/jms/spout/MockJmsProvider.java:50
↓ 1 callers
Method
construct
()
examples/storm-starter/src/jvm/org/apache/storm/starter/ReachTopology.java:72
↓ 1 callers
Method
constructBlobCurrentSymlinkName
(Path baseDir, String key)
storm-server/src/main/java/org/apache/storm/localizer/LocalizedResource.java:131
↓ 1 callers
Method
constructInsertQuery
(String tableName, List<List<Column>> columnLists)
external/storm-jdbc/src/main/java/org/apache/storm/jdbc/common/JdbcClient.java:96
↓ 1 callers
Method
constructVersionFileName
(Path baseDir, String key)
storm-server/src/main/java/org/apache/storm/localizer/LocalizedResource.java:113
↓ 1 callers
Method
consumeImpl
Non blocking. Returns immediately if Q is empty. Returns number of elements consumed from Q.
storm-client/src/jvm/org/apache/storm/utils/JCQueue.java:106
↓ 1 callers
Method
containerFilePath
(String dir)
storm-server/src/main/java/org/apache/storm/utils/ServerUtils.java:268
↓ 1 callers
Method
containsValue
(Object v)
storm-clojure/src/main/java/org/apache/storm/clojure/ClojureTuple.java:405
↓ 1 callers
Method
convertArtifactToJarFileName
(String artifact)
storm-client/src/jvm/org/apache/storm/dependency/DependencyUploader.java:143
↓ 1 callers
Method
convertBase64MapToBinaryMap
(Map<String, String> base64Map)
examples/storm-redis-examples/src/main/java/org/apache/storm/redis/tools/Base64ToBinaryStateMigrationUtil.java:97
↓ 1 callers
Method
convertDuration
(double duration)
storm-client/src/jvm/org/apache/storm/executor/Executor.java:467
↓ 1 callers
Method
convertExecutorBeats
Ensures that we only return heartbeats for executors assigned to this worker.
storm-client/src/jvm/org/apache/storm/cluster/ClusterUtils.java:259
↓ 1 callers
Method
convertExecutorStats
convert thrift ExecutorStats structure into a java HashMap.
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:1398
↓ 1 callers
Method
convertExecutorZkHbs
Convert Long Executor Ids in ZkHbs to Integer ones structure to java maps.
storm-client/src/jvm/org/apache/storm/stats/ClientStatsUtil.java:68
↓ 1 callers
Method
convertExecutorsStats
convert executors stats into a HashMap, note that ExecutorStats are remained unchanged.
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:1383
↓ 1 callers
Method
convertFileSetToFiles
(FileSet source)
storm-buildtools/storm-maven-plugins/src/main/java/org/apache/storm/maven/plugin/versioninfo/VersionInfoMojo.java:80
← previous
next →
6,801–6,900 of 27,770, ranked by callers