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
↓ 2 callers
Method
asString
(byte[] data)
storm-client/src/jvm/org/apache/storm/trident/topology/state/TransactionalState.java:110
↓ 2 callers
Method
asStringList
(Object o)
storm-server/src/main/java/org/apache/storm/daemon/supervisor/BasicContainer.java:415
↓ 2 callers
Method
assertContains
(List<ArtifactResult> results, String groupId, String artifactId, String version)
storm-submit-tools/src/test/java/org/apache/storm/submit/dependency/DependencyResolverTest.java:70
↓ 2 callers
Method
assertLoop
(Predicate<Object> condition, Object... conditionParams)
storm-server/src/test/java/org/apache/storm/AssertLoop.java:32
↓ 2 callers
Method
assertMetrics
(List<String> elements, String endpoint, Map<String, Object> conf)
external/storm-metrics-prometheus/src/test/java/org/apache/storm/metrics/prometheus/PrometheusPreparableReporterTest.java:162
↓ 2 callers
Method
assertNElements
(Collection<?> collection, int expectedCount)
integration-test/src/test/java/org/apache/storm/st/utils/AssertUtil.java:57
↓ 2 callers
Method
assertTopologiesNotBeenEvicted
(Cluster cluster, Class strategyClass, Set<String> evictedTopologies, String... topoNames)
storm-server/src/test/java/org/apache/storm/scheduler/resource/TestUtilsForResourceAwareScheduler.java:536
↓ 2 callers
Method
assign
(WorkerSlot slot, Collection<ExecutorDetails> executors)
storm-server/src/main/java/org/apache/storm/scheduler/SchedulerAssignmentImpl.java:141
↓ 2 callers
Method
assignBoundAckersForNewWorkerSlot
<p> Determine how many bound ackers to put into the given workerSlot. Then try to assign the ackers one by one into this workerSlot upto the calculate
storm-server/src/main/java/org/apache/storm/scheduler/resource/strategies/scheduling/BaseResourceAwareStrategy.java:566
↓ 2 callers
Method
assignCurrentExecutor
Attempt to assign current executor (execIndex points to) to worker and node. Assignment validity check is done before calling this method. @param exe
storm-server/src/main/java/org/apache/storm/scheduler/resource/strategies/scheduling/SchedulingSearcherState.java:214
↓ 2 callers
Method
assignInternal
(WorkerSlot ws, String topId, boolean dontThrow)
storm-server/src/main/java/org/apache/storm/scheduler/multitenant/Node.java:205
↓ 2 callers
Method
assignSingleExecutor
Assign a single executor to a slot, even if other things are in the slot. @param ws the slot to assign it to. @param exec the executor to assign. @par
storm-server/src/main/java/org/apache/storm/scheduler/resource/RasNode.java:334
↓ 2 callers
Method
assignSlotTo
Assign a slot to the given node. @param n the node to assign a slot to. @return true if there are more slots to assign else false.
storm-server/src/main/java/org/apache/storm/scheduler/multitenant/NodePool.java:258
↓ 2 callers
Method
assignSlots
this is called after the assignment is changed in ZK.
storm-server/src/main/java/org/apache/storm/scheduler/INimbus.java:33
↓ 2 callers
Method
assignmentChangedNodes
Diff old/new assignment to find nodes which assigned assignments has changed. @param oldAss old assigned assignment @param newAss new assigned assign
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:1648
↓ 2 callers
Method
authenticated
(Channel c)
storm-client/src/jvm/org/apache/storm/messaging/netty/ISaslServer.java:22
↓ 2 callers
Method
available
()
storm-client/src/jvm/org/apache/storm/blobstore/NimbusBlobStore.java:359
↓ 2 callers
Method
awaitLeadership
Wait for the caller to gain leadership. This should only be used in single-Nimbus clusters, and is only useful to allow testing code to wait for a Loc
storm-client/src/jvm/org/apache/storm/nimbus/ILeaderElector.java:63
↓ 2 callers
Method
backpressurePath
(String stormId, String node, Long port)
storm-client/src/jvm/org/apache/storm/cluster/ClusterUtils.java:167
↓ 2 callers
Method
backpressureTopologies
Get backpressure topologies. Note: In Storm 2.0. Retained for enabling transition from 1.x. Will be removed soon.
storm-client/src/jvm/org/apache/storm/cluster/IStormClusterState.java:156
↓ 2 callers
Method
basicFailureTest
(String confKey, Object confValue, ConstraintSolverStrategy cs, boolean conso
storm-server/src/test/java/org/apache/storm/scheduler/resource/strategies/scheduling/TestConstraintSolverStrategy.java:402
↓ 2 callers
Method
basicSupervisorDetailsMap
(IStormClusterState state)
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:982
↓ 2 callers
Method
batchConf
()
storm-core/test/jvm/org/apache/storm/messaging/netty/NettyTest.java:328
↓ 2 callers
Method
batchRetrieve
(S state, List<TridentTuple> args)
storm-client/src/jvm/org/apache/storm/trident/state/QueryFunction.java:21
↓ 2 callers
Method
blacklistHost
(String host)
storm-server/src/main/java/org/apache/storm/scheduler/Cluster.java:328
↓ 2 callers
Method
blacklistHostsAndSortNodes
( Set<String> blackListedHosts, Collection<SupervisorDetails> sups, Cluster cluster, TopologyDetails t
storm-server/src/test/java/org/apache/storm/scheduler/resource/strategies/scheduling/sorter/TestNodeSorterHostProximity.java:984
↓ 2 callers
Method
blacklistHostsAndSortNodes
( Set<String> blackListedHosts, Collection<SupervisorDetails> sups, Cluster cluster, TopologyDetails t
storm-server/src/test/java/org/apache/storm/scheduler/resource/strategies/scheduling/sorter/TestRoundRobinNodeSorterHostProximity.java:890
↓ 2 callers
Method
blobSync
()
storm-server/src/main/java/org/apache/storm/blobstore/LocalFsBlobStore.java:151
↓ 2 callers
Method
blobstore
(Runnable callback)
storm-client/src/jvm/org/apache/storm/cluster/IStormClusterState.java:220
↓ 2 callers
Method
boltAckedTuple
(String sourceComponentId, String sourceStreamId, long latencyMs)
storm-client/src/jvm/org/apache/storm/metrics2/TaskMetrics.java:72
↓ 2 callers
Method
boltEmit
(String streamId, Collection<Tuple> anchors, List<Object> values, Integer t
storm-client/src/jvm/org/apache/storm/executor/bolt/BoltOutputCollectorImpl.java:82
↓ 2 callers
Method
boltExecute
(BoltExecuteInfo info)
storm-client/src/jvm/org/apache/storm/hooks/ITaskHook.java:35
↓ 2 callers
Method
boltExecuteTuple
(String sourceComponentId, String sourceStreamId, long latencyMs)
storm-client/src/jvm/org/apache/storm/metrics2/TaskMetrics.java:111
↓ 2 callers
Method
boltFailedTuple
(String sourceComponentId, String sourceStreamId)
storm-client/src/jvm/org/apache/storm/metrics2/TaskMetrics.java:90
↓ 2 callers
Method
buffer_for_binary_arg
()
storm-client/src/jvm/org/apache/storm/generated/JavaObjectArg.java:494
↓ 2 callers
Method
buffer_for_custom_serialized
()
storm-client/src/jvm/org/apache/storm/generated/Grouping.java:647
↓ 2 callers
Method
buffer_for_message_blob
()
storm-client/src/jvm/org/apache/storm/generated/HBMessageData.java:513
↓ 2 callers
Method
buffer_for_serialized_java
()
storm-client/src/jvm/org/apache/storm/generated/ComponentObject.java:326
↓ 2 callers
Method
build
(IBackingMap<T> backing)
storm-client/src/jvm/org/apache/storm/trident/state/map/NonTransactionalMap.java:27
↓ 2 callers
Method
build
Builds container for single Redis environment. @param config configuration for JedisPool @return container for single Redis environment
external/storm-redis/src/main/java/org/apache/storm/redis/common/container/JedisCommandsContainerBuilder.java:32
↓ 2 callers
Method
buildDownloadFile
Build a Response object representing download a file. @param contentDispositionName The name to set in the Content-Disposition header @param file fil
storm-webapp/src/main/java/org/apache/storm/daemon/logviewer/utils/LogviewerResponseBuilder.java:77
↓ 2 callers
Method
buildLocalTasksUnevenLoadMapping
(List<Integer> availableTasks)
storm-client/test/jvm/org/apache/storm/grouping/LoadAwareShuffleGroupingTest.java:349
↓ 2 callers
Method
cacheNCheckFields
(RecordTranslator<K, V> translator)
external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/ByTopicRecordTranslator.java:122
↓ 2 callers
Method
calcThroughput
Calculate throughput. @return events/sec
examples/storm-perf/src/main/java/org/apache/storm/perf/ThroughputMeter.java:37
↓ 2 callers
Method
calculateMemoryLimit
(final WorkerResources resources, final int memOnHeap)
storm-server/src/main/java/org/apache/storm/daemon/supervisor/BasicContainer.java:795
↓ 2 callers
Function
calculate_file_hash
(filename)
bin/docker-to-squash.py:258
↓ 2 callers
Method
canSchedule
()
storm-server/src/main/java/org/apache/storm/scheduler/resource/ResourceAwareScheduler.java:334
↓ 2 callers
Method
captureTopology
Rewrites a topology so that all the tuples flowing through it are captured. @param topology the topology to rewrite @return the modified topology and
storm-server/src/main/java/org/apache/storm/Testing.java:341
↓ 2 callers
Method
captureTopology
(Nimbus.Iface client, TopologySummary topologySummary)
examples/storm-loadgen/src/main/java/org/apache/storm/loadgen/CaptureLoad.java:87
↓ 2 callers
Method
checkFailures
()
storm-client/src/jvm/org/apache/storm/windowing/TimeTriggerPolicy.java:93
↓ 2 callers
Method
checkFilesExist
(Collection<File> dependencies)
storm-client/src/jvm/org/apache/storm/dependency/DependencyUploader.java:187
↓ 2 callers
Method
checkFinish
(TrackedBatch tracked, Tuple tuple, TupleType type)
storm-client/src/jvm/org/apache/storm/trident/topology/TridentBoltExecutor.java:149
↓ 2 callers
Method
checkIfStateContainsCurrentNimbusHost
(List<String> stateInfoList, NimbusInfo nimbusInfo)
storm-server/src/main/java/org/apache/storm/blobstore/KeySequenceNumber.java:204
↓ 2 callers
Method
checkInitialization
Checks if the topology's resource requirements are initialized. Will modify topologyResources by adding the appropriate defaults @param topologyResour
examples/storm-loadgen/src/main/java/org/apache/storm/loadgen/CaptureLoad.java:436
↓ 2 callers
Method
checkMatching
(String metricName, List<Pattern> patterns, boolean valueWhenMatched)
storm-client/src/jvm/org/apache/storm/metric/filter/FilterByMetricName.java:89
↓ 2 callers
Method
checkPermission
(String key, Subject who, int mask)
storm-server/src/main/java/org/apache/storm/blobstore/LocalFsBlobStore.java:367
↓ 2 callers
Method
checkTopoPermission
(String principal, String user, Set<String> userGroups, Map<String, Ob
storm-client/src/jvm/org/apache/storm/security/auth/authorizer/SimpleACLAuthorizer.java:177
↓ 2 callers
Method
checkUserGroupAllowed
(Set<String> userGroups, Set<String> configuredGroups)
storm-client/src/jvm/org/apache/storm/security/auth/authorizer/SupervisorSimpleACLAuthorizer.java:144
↓ 2 callers
Method
checkValidReader
(String readerType)
external/storm-hdfs/src/main/java/org/apache/storm/hdfs/spout/HdfsSpout.java:120
↓ 2 callers
Function
check_image_for_magic_file
(magic_file, skopeo_dir, layers)
bin/docker-to-squash.py:678
↓ 2 callers
Method
chooseTaskIndex
(List<T> keys, int numTasks)
storm-client/src/jvm/org/apache/storm/utils/TupleUtils.java:37
↓ 2 callers
Method
chooseTasks
(int taskId, List<Object> values)
storm-client/src/jvm/org/apache/storm/grouping/ShuffleGrouping.java:40
↓ 2 callers
Method
cleanUpForRestart
Clean up the container partly preparing for restart. By default delete all of the temp directories we are going to get a new worker_id anyways. POST C
storm-server/src/main/java/org/apache/storm/daemon/supervisor/Container.java:479
↓ 2 callers
Method
cleanUpTemp
(String baseName)
storm-server/src/main/java/org/apache/storm/localizer/LocallyCachedTopologyBlob.java:263
↓ 2 callers
Method
cleanup
This function will be called when the worker needs to shutdown. This function should include logic to clean up after a worker is shutdown. @param use
storm-server/src/main/java/org/apache/storm/container/ResourceIsolationInterface.java:53
↓ 2 callers
Method
cleanup
()
storm-server/src/main/java/org/apache/storm/metric/api/IClusterMetricsConsumer.java:24
↓ 2 callers
Method
cleanup
(String id)
storm-server/src/main/java/org/apache/storm/daemon/drpc/DRPC.java:146
↓ 2 callers
Method
cleanup
Actually cleanup the blobs to try and get below the target cache size. @param store the blobs store client used to check if the blob has been deleted
storm-server/src/main/java/org/apache/storm/localizer/LocalizedResourceRetentionSet.java:84
↓ 2 callers
Method
cleanup
called once when the system is shutting down, should be idempotent.
external/storm-jdbc/src/main/java/org/apache/storm/jdbc/common/ConnectionProvider.java:36
↓ 2 callers
Method
cleanupAll
(long timeoutMs, DRPCExecutionException exp)
storm-server/src/main/java/org/apache/storm/daemon/drpc/DRPC.java:153
↓ 2 callers
Method
cleanupCutoffAgeMillis
(long nowMillis)
storm-webapp/src/main/java/org/apache/storm/daemon/logviewer/utils/LogCleaner.java:313
↓ 2 callers
Method
clearCredentials
(Subject subject, KerberosTicket tgt)
storm-client/src/jvm/org/apache/storm/security/auth/kerberos/AutoTGT.java:81
↓ 2 callers
Method
close
()
storm-client/src/jvm/org/apache/storm/LogWriter.java:79
↓ 2 callers
Method
close
()
storm-client/src/jvm/org/apache/storm/utils/JCQueue.java:69
↓ 2 callers
Method
close
()
storm-server/src/main/java/org/apache/storm/LocalCluster.java:620
↓ 2 callers
Method
close
()
storm-server/src/main/java/org/apache/storm/metric/timed/TimerDecorated.java:34
↓ 2 callers
Method
close
()
storm-server/src/main/java/org/apache/storm/daemon/supervisor/Supervisor.java:501
↓ 2 callers
Method
closeChannelAndReconnect
Schedule a reconnect if we closed a non-null channel, and acquired the right to provide a replacement by successfully setting a null to the channel fi
storm-client/src/jvm/org/apache/storm/messaging/netty/Client.java:451
↓ 2 callers
Method
closeOutputFile
()
external/storm-hdfs/src/main/java/org/apache/storm/hdfs/trident/HdfsState.java:205
↓ 2 callers
Method
closeReaderAndResetTrackers
()
external/storm-hdfs/src/main/java/org/apache/storm/hdfs/spout/HdfsSpout.java:367
↓ 2 callers
Method
combine
(T val1, T val2)
storm-client/src/jvm/org/apache/storm/trident/operation/CombinerAggregator.java:23
↓ 2 callers
Method
combine
(List<Measurements> measurements, Integer start, Integer count)
examples/storm-loadgen/src/main/java/org/apache/storm/loadgen/LoadMetricsServer.java:288
↓ 2 callers
Method
commitIfNecessary
()
external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaSpout.java:660
↓ 2 callers
Method
compare
(Node o1, Node o2)
storm-server/src/main/java/org/apache/storm/scheduler/multitenant/Node.java:39
↓ 2 callers
Method
compareTo
(Rankable other)
examples/storm-starter/src/jvm/org/apache/storm/starter/tools/RankableObjectWithFields.java:84
↓ 2 callers
Method
compareTo
(JmsMessageID jmsMessageID)
external/storm-jms/src/main/java/org/apache/storm/jms/spout/JmsMessageID.java:34
↓ 2 callers
Method
complete
()
storm-server/src/main/java/org/apache/storm/localizer/PortAndAssignment.java:29
↓ 2 callers
Method
componentBoltSubscriptions
(Component component)
storm-client/src/jvm/org/apache/storm/coordination/BatchSubtopologyBuilder.java:123
↓ 2 callers
Method
componentToExecs
(String comp)
storm-server/src/main/java/org/apache/storm/scheduler/TopologyDetails.java:188
↓ 2 callers
Method
computeAggCapacity
(Map m, Integer uptime)
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:2126
↓ 2 callers
Method
computeExecutors
(StormBase base, Map<String, Object> topoConf, StormTopology
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:2107
↓ 2 callers
Method
computeScheduledTopologyMemory
(Cluster cluster, TopologyDetails td)
storm-server/src/main/java/org/apache/storm/scheduler/resource/ResourceAwareScheduler.java:400
↓ 2 callers
Method
computeTotalCount
(T obj)
examples/storm-starter/src/jvm/org/apache/storm/starter/tools/SlotBasedCounter.java:67
↓ 2 callers
Method
config
This function returns the scheduler's configuration. @return The scheduler's configuration.
storm-server/src/main/java/org/apache/storm/scheduler/IScheduler.java:41
↓ 2 callers
Method
constructEmptyStormTopology
()
storm-server/src/test/java/org/apache/storm/localizer/AsyncLocalizerTest.java:341
↓ 2 callers
Method
constructExpectedArchivesDir
(String base, String user)
storm-server/src/test/java/org/apache/storm/localizer/AsyncLocalizerTest.java:361
↓ 2 callers
Method
constructExpectedFilesDir
(String base, String user)
storm-server/src/test/java/org/apache/storm/localizer/AsyncLocalizerTest.java:357
↓ 2 callers
Method
containerPidFile
(String workerId)
storm-server/src/main/java/org/apache/storm/container/oci/RuncLibContainerManager.java:160
↓ 2 callers
Method
convertExecutorBeats
convert thrift executor heartbeats into a java HashMap.
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:1325
↓ 2 callers
Method
convertPatternStringsToPatternInstances
(List<String> patterns)
storm-client/src/jvm/org/apache/storm/metric/filter/FilterByMetricName.java:75
← previous
next →
4,201–4,300 of 27,770, ranked by callers