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
convertSpecificStats
(SpoutStats stats)
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:1416
↓ 2 callers
Method
convertToArray
(Map<Integer, V> srcMap, int start)
storm-client/src/jvm/org/apache/storm/utils/Utils.java:1761
↓ 2 callers
Method
convertZkWorkerHb
convert a thrift worker heartbeat into a java HashMap.
storm-server/src/main/java/org/apache/storm/stats/StatsUtil.java:1369
↓ 2 callers
Method
coordinatorPath
(Configuration configuration, String txid)
external/storm-kafka-migration/src/main/java/org/apache/storm/kafka/migration/KafkaTridentSpoutMigration.java:172
↓ 2 callers
Method
copy
Creates a (defensive) copy of itself.
examples/storm-starter/src/jvm/org/apache/storm/starter/tools/Rankings.java:147
↓ 2 callers
Method
copy
Must be able to copy the rotation policy.
external/storm-hdfs/src/main/java/org/apache/storm/hdfs/bolt/rotation/FileRotationPolicy.java:47
↓ 2 callers
Method
copyDirectory
Copy a directory. @param fromDir from where @param toDir to where @throws IOException on any error
storm-client/src/jvm/org/apache/storm/daemon/supervisor/IAdvancedFSOps.java:72
↓ 2 callers
Method
copyInputStreamToBlobOutputStream
(InputStream is, AtomicOutputStream os)
storm-core/src/jvm/org/apache/storm/command/Blobstore.java:297
↓ 2 callers
Method
couldEverFit
Is there any possibility that exec could ever fit on this node. @param exec the executor to schedule @param td the topology the executor is a part of
storm-server/src/main/java/org/apache/storm/scheduler/resource/RasNode.java:413
↓ 2 callers
Method
couldFit
Is there any possibility that a resource request could ever fit on this. @param minWorkerCpu the configured minimum worker CPU @param requestedResourc
storm-server/src/main/java/org/apache/storm/scheduler/resource/normalization/NormalizedResourceOffer.java:199
↓ 2 callers
Method
count
(int[] choices, List<Integer>[] rets)
storm-client/test/jvm/org/apache/storm/grouping/LoadAwareShuffleGroupingTest.java:125
↓ 2 callers
Method
countNonZeroLengthFiles
(String path)
external/storm-hdfs/src/test/java/org/apache/storm/hdfs/bolt/TestSequenceFileBolt.java:168
↓ 2 callers
Method
createBackPressureWaitStrategy
(Map<String, Object> topologyConf)
storm-client/src/jvm/org/apache/storm/policy/IWaitStrategy.java:27
↓ 2 callers
Method
createBeatBoltStats
Utility method for creating a template for Bolt stats. @return Empty template map for Bolt statistics.
storm-core/test/jvm/org/apache/storm/stats/TestStatsUtil.java:119
↓ 2 callers
Method
createBlobFromStream
(final String key, final InputStream is, final SettableBlobMeta meta)
storm-core/src/jvm/org/apache/storm/command/Blobstore.java:283
↓ 2 callers
Method
createCSSClusterConfig
(double compPcore, double compOnHeap, double compOffHeap, Map<
storm-server/src/test/java/org/apache/storm/scheduler/resource/TestUtilsForResourceAwareScheduler.java:110
↓ 2 callers
Method
createConf
()
storm-client/test/jvm/org/apache/storm/grouping/LoadAwareShuffleGroupingTest.java:54
↓ 2 callers
Method
createConfigLoader
The user interface to create an IConfigLoader instance. It iterates all the implementations of IConfigLoaderFactory and finds the one which supports t
storm-server/src/main/java/org/apache/storm/scheduler/utils/ConfigLoaderFactoryService.java:43
↓ 2 callers
Method
createEmittingContext
(ProcessorNode processorNode)
storm-client/src/jvm/org/apache/storm/streams/ProcessorBoltDelegate.java:243
↓ 2 callers
Method
createFetchedOffsetsMetadata
(Set<TopicPartition> assignedPartitions)
external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaSpout.java:497
↓ 2 callers
Method
createFile
(String fileName)
examples/storm-starter/src/jvm/org/apache/storm/starter/BlobStoreAPIWordCountTopology.java:145
↓ 2 callers
Method
createFileReader
Creates a reader that reads from beginning of file. @param file file to read
external/storm-hdfs/src/main/java/org/apache/storm/hdfs/spout/HdfsSpout.java:687
↓ 2 callers
Method
createGrasClusterConfig
(double compPcore, double compOnHeap, double compOffHeap, Map
storm-server/src/test/java/org/apache/storm/scheduler/resource/TestUtilsForResourceAwareScheduler.java:130
↓ 2 callers
Method
createHostToSupervisorMap
(final List<String> blacklistedNodeIds, Cluster cluster)
storm-server/src/main/java/org/apache/storm/scheduler/blacklist/strategies/DefaultBlacklistStrategy.java:197
↓ 2 callers
Method
createKeyManager
(KeyStore keystore, String keystorePassword)
storm-client/src/jvm/org/apache/storm/security/auth/tls/ReloadableX509KeyManager.java:59
↓ 2 callers
Method
createMessage
(String fileName)
storm-client/src/jvm/org/apache/storm/dependency/FileNotAvailableException.java:24
↓ 2 callers
Method
createNewWorkerId
Create a new worker ID for this process and store in in this object and in the local state. Never call this if a worker is currently up and running.
storm-server/src/main/java/org/apache/storm/daemon/supervisor/BasicContainer.java:212
↓ 2 callers
Method
createNimbusClient
(Map<String, Object> conf, String asUser, Integer timeout)
storm-client/src/jvm/org/apache/storm/utils/NimbusClient.java:235
↓ 2 callers
Method
createNoWaitRetryService
()
external/storm-kafka-client/src/test/java/org/apache/storm/kafka/spout/KafkaSpoutRetryExponentialBackoffTest.java:37
↓ 2 callers
Method
createOutputFile
()
external/storm-hdfs/src/main/java/org/apache/storm/hdfs/trident/HdfsState.java:207
↓ 2 callers
Method
createPassword
Compute HMAC of the identifier using the secret key and return the output as password. @param identifier the bytes of the identifier @param key
storm-client/src/jvm/org/apache/storm/security/auth/workertoken/WorkerTokenSigner.java:54
↓ 2 callers
Method
createQueue
(String name, int queueSize)
storm-client/test/jvm/org/apache/storm/utils/JCQueueBackpressureTest.java:27
↓ 2 callers
Method
createRemoteTopology
()
storm-client/src/jvm/org/apache/storm/drpc/LinearDRPCTopologyBuilder.java:93
↓ 2 callers
Method
createSaslServer
(String mechanism, String protocol, String serverName, Map<String,
storm-client/src/jvm/org/apache/storm/security/auth/plain/SaslPlainServer.java:141
↓ 2 callers
Method
createSslContext
(Map<String, Object> topoConf, boolean forServer)
storm-client/src/jvm/org/apache/storm/messaging/netty/NettyTlsUtils.java:38
↓ 2 callers
Method
createSupervisorClient
()
storm-client/src/jvm/org/apache/storm/utils/SupervisorClient.java:65
↓ 2 callers
Method
createTestBlob
(String testKey, SettableBlobMeta meta)
storm-client/test/jvm/org/apache/storm/blobstore/ClientBlobStoreTest.java:101
↓ 2 callers
Method
createTopoDetailsArray
Create an array of TopologyDetails by reading serialized files for topology and configuration in the resource path. Skip topologies with no executors/
storm-server/src/test/java/org/apache/storm/scheduler/resource/strategies/scheduling/TestLargeCluster.java:183
↓ 2 callers
Method
createTopology
(DRPCSpout spout)
storm-client/src/jvm/org/apache/storm/drpc/LinearDRPCTopologyBuilder.java:97
↓ 2 callers
Method
createTp
(int partition)
external/storm-kafka-client/src/test/java/org/apache/storm/kafka/spout/subscription/RoundRobinManualPartitionerTest.java:36
↓ 2 callers
Method
createTrustManager
(String trustStorePath, String keystorePassword)
storm-client/src/jvm/org/apache/storm/security/auth/tls/ReloadableX509TrustManager.java:81
↓ 2 callers
Method
create_sequential
Path will be appended with a monotonically increasing integer, a new node will be created there, and data will be put at that node. @param path The p
storm-client/src/jvm/org/apache/storm/cluster/IStateStorage.java:59
↓ 2 callers
Method
created_at
(self)
dev-tools/github/__init__.py:73
↓ 2 callers
Method
customGrouping
(final String component, final CustomStreamGrouping grouping)
storm-client/src/jvm/org/apache/storm/trident/topology/TridentTopologyBuilder.java:699
↓ 2 callers
Method
customGrouping
(String componentId, CustomStreamGrouping grouping)
storm-client/src/jvm/org/apache/storm/topology/TopologyBuilder.java:745
↓ 2 callers
Method
customGrouping
(final CustomStreamGrouping grouping)
storm-client/src/jvm/org/apache/storm/drpc/LinearDRPCTopologyBuilder.java:365
↓ 2 callers
Method
customGrouping
(CustomStreamGrouping grouping)
storm-client/src/jvm/org/apache/storm/drpc/LinearDRPCInputDeclarer.java:53
↓ 2 callers
Method
customGrouping
(final String component, final CustomStreamGrouping grouping)
storm-client/src/jvm/org/apache/storm/coordination/BatchSubtopologyBuilder.java:391
↓ 2 callers
Method
customResourceMapEquality
This method compares Resource Maps while considering any resources are NULL to be 0.0 @param firstMap Resource Map A @param secondMap Resource Map B
storm-server/src/main/java/org/apache/storm/utils/EquivalenceUtils.java:117
↓ 2 callers
Method
deactivate
(java.lang.String name)
storm-client/src/jvm/org/apache/storm/generated/Nimbus.java:42
↓ 2 callers
Method
debug
(String topoName, String componentId, boolean enable, double samplingPercentage)
storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java:3658
↓ 2 callers
Method
declare
(InputDeclarer declarer)
storm-client/src/jvm/org/apache/storm/coordination/BatchSubtopologyBuilder.java:132
↓ 2 callers
Method
declareCheckpointStream
(OutputFieldsDeclarer declarer)
storm-client/src/jvm/org/apache/storm/topology/BaseStatefulBoltExecutor.java:126
↓ 2 callers
Method
declareOutputFields
declare what are the fields that this code will output. @param declarer OutputFieldsDeclarer
external/storm-redis/src/main/java/org/apache/storm/redis/common/mapper/RedisFilterMapper.java:26
↓ 2 callers
Method
decodeKey
Decode key. @param encodedKey the value of key (KRAW type) @return the decoded value of key (K type)
storm-client/src/jvm/org/apache/storm/state/StateEncoder.java:41
↓ 2 callers
Method
defaultSchedule
(Topologies topologies, Cluster cluster)
storm-server/src/main/java/org/apache/storm/scheduler/DefaultScheduler.java:74
↓ 2 callers
Method
delete
()
external/storm-hdfs-blobstore/src/main/java/org/apache/storm/hdfs/blobstore/HdfsBlobStoreFile.java:91
↓ 2 callers
Method
deleteDir
(String dir)
storm-client/src/jvm/org/apache/storm/container/cgroup/CgroupUtils.java:36
↓ 2 callers
Method
deleteKeyIgnoringFileNotFound
(String key)
storm-server/src/main/java/org/apache/storm/blobstore/LocalFsBlobStore.java:373
↓ 2 callers
Method
deleteNode
(CuratorFramework zk, String path)
storm-client/src/jvm/org/apache/storm/zookeeper/ClientZookeeper.java:152
↓ 2 callers
Method
deleteOldestWhileTooLarge
If totalSize of files exceeds the either the per-worker quota or global quota, Logviewer deletes oldest inactive log files in a worker directory or in
storm-webapp/src/main/java/org/apache/storm/daemon/logviewer/utils/DirectoryCleaner.java:90
↓ 2 callers
Method
deletePartition
(long pid)
storm-client/src/jvm/org/apache/storm/windowing/persistence/WindowState.java:261
↓ 2 callers
Method
deletePulseId
(String path)
storm-server/src/main/java/org/apache/storm/pacemaker/Pacemaker.java:198
↓ 2 callers
Method
deleteVersion
(long version)
storm-client/src/jvm/org/apache/storm/utils/VersionedStore.java:105
↓ 2 callers
Method
delete_node_blobstore
Allows us to delete the znodes within /storm/blobstore/key_name whose znodes start with the corresponding nimbusHostPortInfo. @param path
storm-client/src/jvm/org/apache/storm/cluster/IStateStorage.java:209
↓ 2 callers
Method
deriveNumWindowChunksFrom
(int windowLengthInSeconds, int windowUpdateFrequencyInSeconds)
examples/storm-starter/src/jvm/org/apache/storm/starter/bolt/RollingCountBolt.java:78
↓ 2 callers
Method
deserialize
(ThriftSerializedObject obj, TDeserializer td)
storm-client/src/jvm/org/apache/storm/utils/LocalState.java:84
↓ 2 callers
Method
deserializeKerberosTicket
(final byte[] tgtBytes)
storm-client/src/jvm/org/apache/storm/security/auth/ClientAuthUtils.java:528
↓ 2 callers
Method
detachOnRun
Add -d option. @return the self
storm-server/src/main/java/org/apache/storm/container/docker/DockerRunCommand.java:58
↓ 2 callers
Method
direct
(NullStruct value)
storm-client/src/jvm/org/apache/storm/generated/Grouping.java:186
↓ 2 callers
Method
doExecute
(Tuple tuple)
storm-client/src/jvm/org/apache/storm/topology/StatefulBoltExecutor.java:144
↓ 2 callers
Method
doFilterNullTupleTest
(KafkaSpoutConfig.ProcessingGuarantee processingGuarantee)
external/storm-kafka-client/src/test/java/org/apache/storm/kafka/spout/KafkaSpoutMessagingGuaranteeTest.java:250
↓ 2 callers
Method
doGet
(int port, String func, String args)
storm-webapp/src/test/java/org/apache/storm/daemon/drpc/DRPCServerTest.java:135
↓ 2 callers
Method
doGetCredentials
(CredentialKeyProvider provider, Map<String, String> credentials, String configKey)
external/storm-autocreds/src/main/java/org/apache/storm/common/HadoopCredentialUtil.java:63
↓ 2 callers
Method
doHeartBeat
()
storm-client/src/jvm/org/apache/storm/daemon/worker/Worker.java:388
↓ 2 callers
Method
doPrepare
(Map<String, Object> topoConf, TopologyContext context, OutputCollector collector,
storm-client/src/jvm/org/apache/storm/topology/WindowedBoltExecutor.java:291
↓ 2 callers
Method
doProjection
(ArrayList<Tuple> tuples, FieldSelector[] projectionFields)
storm-client/src/jvm/org/apache/storm/bolt/JoinBolt.java:317
↓ 2 callers
Method
doReactivationTest
(FirstPollOffsetStrategy firstPollOffsetStrategy)
external/storm-kafka-client/src/test/java/org/apache/storm/kafka/spout/KafkaSpoutReactivationTest.java:104
↓ 2 callers
Method
doRequiredTopoFilesExist
(Map<String, Object> conf, String stormId)
storm-client/src/jvm/org/apache/storm/daemon/supervisor/ClientSupervisorUtils.java:43
↓ 2 callers
Method
doRequiredTopoFilesExist
Sanity check if everything the topology needs is there for it to run. @param conf the config of the supervisor @param topologyId the ID of the
storm-client/src/jvm/org/apache/storm/daemon/supervisor/IAdvancedFSOps.java:128
↓ 2 callers
Method
doRotationAndRemoveAllWriters
()
external/storm-hdfs/src/main/java/org/apache/storm/hdfs/bolt/AbstractHdfsBolt.java:255
↓ 2 callers
Method
doSanityCheck
()
storm-client/src/jvm/org/apache/storm/task/GeneralTopologyContext.java:214
↓ 2 callers
Method
doTestBasic
(Map<String, Object> stormConf)
storm-core/test/jvm/org/apache/storm/messaging/netty/NettyTest.java:91
↓ 2 callers
Method
doTestBatch
(Map<String, Object> stormConf)
storm-core/test/jvm/org/apache/storm/messaging/netty/NettyTest.java:296
↓ 2 callers
Method
doTestLargeMessage
(Map<String, Object> stormConf)
storm-core/test/jvm/org/apache/storm/messaging/netty/NettyTest.java:208
↓ 2 callers
Method
doTestLoad
(Map<String, Object> stormConf)
storm-core/test/jvm/org/apache/storm/messaging/netty/NettyTest.java:153
↓ 2 callers
Method
doTestModeCannotReplayTuples
(KafkaSpoutConfig<String, String> spoutConfig)
external/storm-kafka-client/src/test/java/org/apache/storm/kafka/spout/KafkaSpoutMessagingGuaranteeTest.java:141
↓ 2 callers
Method
doTestModeDisregardsMaxUncommittedOffsets
(KafkaSpoutConfig<String, String> spoutConfig)
external/storm-kafka-client/src/test/java/org/apache/storm/kafka/spout/KafkaSpoutMessagingGuaranteeTest.java:106
↓ 2 callers
Method
doTestServerAlwaysReconnects
(Map<String, Object> stormConf)
storm-core/test/jvm/org/apache/storm/messaging/netty/NettyTest.java:344
↓ 2 callers
Method
doTestServerDelayed
(Map<String, Object> stormConf)
storm-core/test/jvm/org/apache/storm/messaging/netty/NettyTest.java:247
↓ 2 callers
Method
dockerCidFilePath
(String workerId)
storm-server/src/main/java/org/apache/storm/container/docker/DockerManager.java:336
↓ 2 callers
Function
docker_to_squash
(layer_dir, layer, working_dir)
bin/docker-to-squash.py:659
↓ 2 callers
Function
does_image_have_dead_perms
(image)
bin/docker-to-squash.py:1453
↓ 2 callers
Method
downgrade
(LocalityScope current)
storm-client/src/jvm/org/apache/storm/grouping/LoadAwareShuffleGrouping.java:324
↓ 2 callers
Method
downloadFile
Checks authorization for the log file and download. @param host host address @param fileName file to download @param user username @param isDaemon tr
storm-webapp/src/main/java/org/apache/storm/daemon/logviewer/utils/LogFileDownloader.java:67
↓ 2 callers
Method
drain
()
storm-client/src/jvm/org/apache/storm/messaging/netty/MessageBuffer.java:44
↓ 2 callers
Method
drainAllChangingBlobs
Drop all of the changingBlobs and pendingChangingBlobs. <p>PRECONDITION: container is null @param dynamicState current state. @return the next state
storm-server/src/main/java/org/apache/storm/daemon/supervisor/Slot.java:320
↓ 2 callers
Method
drop
Drop the first N elements and create a new list. @param list the list @param count element count to drop @return newly created sublist that drops the
storm-webapp/src/main/java/org/apache/storm/daemon/utils/ListFunctionalSupport.java:87
↓ 2 callers
Method
dropMessages
(Iterator<TaskMessage> msgs)
storm-client/src/jvm/org/apache/storm/messaging/netty/Client.java:393
↓ 2 callers
Method
dumpState
(PrintStream stream)
external/storm-hdfs/src/main/java/org/apache/storm/hdfs/spout/ProgressTracker.java:57
← previous
next →
4,301–4,400 of 27,770, ranked by callers