Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/TALKDATA/JavaBigData
/ functions
Functions
982 in github.com/TALKDATA/JavaBigData
⨍
Functions
982
◇
Types & classes
173
↓ 5 callers
Method
createIndexRequest
@param client ElasticSearch {@link Client} to prepare index from @param indexPrefix Prefix of index name to use -- as configured on
code/flume-ng-sinks/flume-ng-elasticsearch-sink/src/main/java/org/apache/flume/sink/elasticsearch/ElasticSearchIndexRequestBuilderFactory.java:55
↓ 5 callers
Method
getContextForRetryTests
()
code/flume-ng-sinks/flume-hdfs-sink/src/test/java/org/apache/flume/sink/hdfs/TestHDFSEventSink.java:1423
↓ 5 callers
Method
getMethod
()
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/test/java/org/apache/flume/sink/solr/morphline/FlumeHttpServletRequestWrapper.java:87
↓ 5 callers
Method
getNextMessageFromConsumer
(String topic)
code/flume-ng-sinks/flume-ng-kafka-sink/src/test/java/org/apache/flume/sink/kafka/util/TestUtil.java:130
↓ 5 callers
Method
getOpenedFilePath
()
code/flume-ng-sinks/flume-hdfs-sink/src/test/java/org/apache/flume/sink/hdfs/MockHDFSWriter.java:51
↓ 5 callers
Method
getRemainingTxns
()
code/flume-ng-sinks/flume-hive-sink/src/main/java/org/apache/flume/sink/hive/HiveWriter.java:121
↓ 5 callers
Method
hflushOrSync
If hflush is available in this version of HDFS, then this method calls hflush, else it calls sync. @param os - The stream to flush/sync @throws IOExce
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/AbstractHDFSWriter.java:261
↓ 5 callers
Method
queryResultSetSize
(String query)
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/test/java/org/apache/flume/sink/solr/morphline/TestMorphlineSolrSink.java:416
↓ 5 callers
Method
reset
()
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/main/java/org/apache/flume/sink/solr/morphline/BlobDeserializer.java:127
↓ 5 callers
Method
startSink
(HiveSink sink, Context context)
code/flume-ng-sinks/flume-hive-sink/src/test/java/org/apache/flume/sink/hive/TestHiveSink.java:408
↓ 5 callers
Method
writeEvents
(HiveWriter writer, int count)
code/flume-ng-sinks/flume-hive-sink/src/test/java/org/apache/flume/sink/hive/TestHiveWriter.java:343
↓ 4 callers
Method
assertEventBodyEquals
(String expected, Event event)
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/test/java/org/apache/flume/sink/solr/morphline/TestBlobDeserializer.java:98
↓ 4 callers
Method
assertMatchAllQuery
(int expectedHits, Event... events)
code/flume-ng-sinks/flume-ng-elasticsearch-sink/src/test/java/org/apache/flume/sink/elasticsearch/AbstractElasticSearchSinkTest.java:122
↓ 4 callers
Method
checkAndThrowInterruptedException
This method if the current thread has been interrupted and throws an exception. @throws InterruptedException
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/BucketWriter.java:646
↓ 4 callers
Method
configure
(Context context)
code/flume-ng-sinks/flume-hive-sink/src/main/java/org/apache/flume/sink/hive/HiveDelimitedTextSerializer.java:69
↓ 4 callers
Method
createDbAndTable
(Driver driver, String databaseName, String tableName, List<String> part
code/flume-ng-sinks/flume-hive-sink/src/test/java/org/apache/flume/sink/hive/TestUtil.java:63
↓ 4 callers
Method
createWriter
Create a new writer. This method also re-loads the dataset so updates to the configuration or a dataset created after Flume starts will be loaded. @
code/flume-ng-sinks/flume-dataset-sink/src/main/java/org/apache/flume/sink/kite/DatasetSink.java:392
↓ 4 callers
Method
doPartitionHeader
This method tests both the default behavior (usePartitionHeader=false) and the behaviour when the partitionId setting is used. Under the default behav
code/flume-ng-sinks/flume-ng-kafka-sink/src/test/java/org/apache/flume/sink/kafka/TestKafkaSink.java:420
↓ 4 callers
Method
ensureOpen
()
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/main/java/org/apache/flume/sink/solr/morphline/BlobDeserializer.java:142
↓ 4 callers
Method
getClient
@param clientType String representation of client type @param hostNames Array of strings that represents hostnames with ports (hostname:port) @p
code/flume-ng-sinks/flume-ng-elasticsearch-sink/src/main/java/org/apache/flume/sink/elasticsearch/client/ElasticSearchClientFactory.java:44
↓ 4 callers
Method
getCodec
(String codecName)
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/HDFSEventSink.java:310
↓ 4 callers
Method
getFilesClosed
()
code/flume-ng-sinks/flume-hdfs-sink/src/test/java/org/apache/flume/sink/hdfs/MockHDFSWriter.java:39
↓ 4 callers
Method
getIndexType
()
code/flume-ng-sinks/flume-ng-elasticsearch-sink/src/main/java/org/apache/flume/sink/elasticsearch/ElasticSearchSink.java:154
↓ 4 callers
Method
getKafkaConsumer
()
code/flume-ng-sinks/flume-ng-kafka-sink/src/test/java/org/apache/flume/sink/kafka/util/TestUtil.java:117
↓ 4 callers
Method
getNameNodeURL
(MiniDFSCluster cluster)
code/flume-ng-sinks/flume-hdfs-sink/src/test/java/org/apache/flume/sink/hdfs/TestHDFSEventSinkOnMiniCluster.java:75
↓ 4 callers
Method
getRowKey
Returns a row-key with the following format: [time in millis]-[random key]-[nonce]
code/flume-ng-sinks/flume-ng-hbase-sink/src/main/java/org/apache/flume/sink/hbase/RegexHbaseEventSerializer.java:144
↓ 4 callers
Method
getSerializer
(String formatType, Context context)
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/SequenceFileSerializerFactory.java:36
↓ 4 callers
Method
getSfWriters
()
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/HDFSEventSink.java:183
↓ 4 callers
Method
getTryCount
()
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/HDFSEventSink.java:555
↓ 4 callers
Method
getZkUrl
()
code/flume-ng-sinks/flume-ng-kafka-sink/src/test/java/org/apache/flume/sink/kafka/util/TestUtil.java:168
↓ 4 callers
Method
handle
Handle a non-recoverable event. @param event The event @param cause The cause of the failure @throws EventDeliveryException The policy failed to hand
code/flume-ng-sinks/flume-dataset-sink/src/main/java/org/apache/flume/sink/kite/policy/FailurePolicy.java:62
↓ 4 callers
Method
handleTransactionFailure
(Transaction txn)
code/flume-ng-sinks/flume-ng-hbase-sink/src/main/java/org/apache/flume/sink/hbase/AsyncHBaseSink.java:556
↓ 4 callers
Method
mark
()
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/main/java/org/apache/flume/sink/solr/morphline/BlobDeserializer.java:121
↓ 4 callers
Method
run
(String filePath)
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/HDFSEventSink.java:63
↓ 4 callers
Method
runDDL
(Driver driver, String sql)
code/flume-ng-sinks/flume-hive-sink/src/test/java/org/apache/flume/sink/hive/TestUtil.java:217
↓ 4 callers
Method
schema
Get the schema from the event headers. @param event The Flume event @return The schema for the event @throws EventDeliveryException A recoverable err
code/flume-ng-sinks/flume-dataset-sink/src/main/java/org/apache/flume/sink/kite/parser/AvroParser.java:175
↓ 4 callers
Method
setClock
(Clock clock)
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/BucketWriter.java:637
↓ 4 callers
Method
stop
()
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/HDFSEventSink.java:475
↓ 3 callers
Method
abort
Aborts the current Txn @throws InterruptedException
code/flume-ng-sinks/flume-hive-sink/src/main/java/org/apache/flume/sink/hive/HiveWriter.java:222
↓ 3 callers
Method
afterCreate
()
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/AvroEventSerializer.java:100
↓ 3 callers
Method
assertBodyQuery
(int expectedHits, Event... events)
code/flume-ng-sinks/flume-ng-elasticsearch-sink/src/test/java/org/apache/flume/sink/elasticsearch/AbstractElasticSearchSinkTest.java:127
↓ 3 callers
Method
beforeClose
()
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/AvroEventSerializer.java:190
↓ 3 callers
Method
build
(Context context)
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/test/java/org/apache/flume/sink/solr/morphline/TestUUIDInterceptor.java:57
↓ 3 callers
Method
build
()
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/main/java/org/apache/flume/sink/solr/morphline/UUIDInterceptor.java:103
↓ 3 callers
Method
configure
(Context context)
code/flume-ng-sinks/flume-hive-sink/src/main/java/org/apache/flume/sink/hive/HiveSink.java:103
↓ 3 callers
Method
configure
(Context context)
code/flume-ng-sinks/flume-dataset-sink/src/main/java/org/apache/flume/sink/kite/DatasetSink.java:160
↓ 3 callers
Method
configure
(Context context)
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/HDFSCompressedDataStream.java:55
↓ 3 callers
Method
configure
(Context context)
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/main/java/org/apache/flume/sink/solr/morphline/BlobHandler.java:65
↓ 3 callers
Method
createAndConfigureMemoryChannel
(HBaseSink sink)
code/flume-ng-sinks/flume-ng-hbase-sink/src/test/java/org/apache/flume/sink/hbase/TestHBaseSink.java:723
↓ 3 callers
Method
destroyConnection
()
code/flume-ng-sinks/flume-irc-sink/src/main/java/org/apache/flume/sink/irc/IRCSink.java:170
↓ 3 callers
Method
doPartitionErrors
This function tests three scenarios: 1. PartitionOption.VALIDBUTOUTOFRANGE: An integer partition is provided, however it exceeds the number of part
code/flume-ng-sinks/flume-ng-kafka-sink/src/test/java/org/apache/flume/sink/kafka/TestKafkaSink.java:362
↓ 3 callers
Method
getActions
Get the actions that should be written out to hbase as a result of this event. This list is written to hbase using the HBase batch API. @return List o
code/flume-ng-sinks/flume-ng-hbase-sink/src/main/java/org/apache/flume/sink/hbase/HbaseEventSerializer.java:53
↓ 3 callers
Method
getAllFiles
(String input)
code/flume-ng-sinks/flume-hdfs-sink/src/test/java/org/apache/flume/sink/hdfs/TestHDFSEventSink.java:750
↓ 3 callers
Method
getClusterName
()
code/flume-ng-sinks/flume-ng-elasticsearch-sink/src/main/java/org/apache/flume/sink/elasticsearch/ElasticSearchSink.java:144
↓ 3 callers
Method
getHeaders
(HttpServletRequest request)
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/main/java/org/apache/flume/sink/solr/morphline/BlobHandler.java:108
↓ 3 callers
Method
getLastUsed
()
code/flume-ng-sinks/flume-hive-sink/src/main/java/org/apache/flume/sink/hive/HiveWriter.java:460
↓ 3 callers
Method
getOpenFileDescriptorCount
()
code/flume-ng-sinks/flume-ng-hbase-sink/src/test/java/org/apache/flume/sink/hbase/TestAsyncHBaseSink.java:453
↓ 3 callers
Method
getWriter
()
code/flume-ng-sinks/flume-dataset-sink/src/main/java/org/apache/flume/sink/kite/DatasetSink.java:231
↓ 3 callers
Method
initContextForIncrementHBaseSerializer
Set up {@link Context} for use with {@link IncrementHBaseSerializer}.
code/flume-ng-sinks/flume-ng-hbase-sink/src/test/java/org/apache/flume/sink/hbase/TestHBaseSink.java:123
↓ 3 callers
Method
open
(String filePath)
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/HDFSCompressedDataStream.java:68
↓ 3 callers
Method
read
()
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/test/java/org/apache/flume/sink/solr/morphline/FlumeHttpServletRequestWrapper.java:45
↓ 3 callers
Method
registerCurrentStream
(FSDataOutputStream outputStream, FileSystem fs, Path destPath)
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/AbstractHDFSWriter.java:106
↓ 3 callers
Method
returnToPool
(LocalMorphlineInterceptor interceptor)
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/main/java/org/apache/flume/sink/solr/morphline/MorphlineInterceptor.java:87
↓ 3 callers
Method
roll
Causes the sink to roll at the next {@link #process()} call.
code/flume-ng-sinks/flume-dataset-sink/src/main/java/org/apache/flume/sink/kite/DatasetSink.java:226
↓ 3 callers
Method
serialize
(Object datum, Schema schema)
code/flume-ng-sinks/flume-dataset-sink/src/test/java/org/apache/flume/sink/kite/TestDatasetSink.java:988
↓ 3 callers
Method
setWriter
(DatasetWriter<GenericRecord> writer)
code/flume-ng-sinks/flume-dataset-sink/src/main/java/org/apache/flume/sink/kite/DatasetSink.java:236
↓ 3 callers
Method
shutdownHBaseClient
()
code/flume-ng-sinks/flume-ng-hbase-sink/src/main/java/org/apache/flume/sink/hbase/AsyncHBaseSink.java:526
↓ 3 callers
Method
start
()
code/flume-ng-sinks/flume-hive-sink/src/main/java/org/apache/flume/sink/hive/HiveSink.java:489
↓ 3 callers
Method
sync
Ensure any handled events are on stable storage. This allows the policy implementation to sync any data that it may not have fully handled. See {@li
code/flume-ng-sinks/flume-dataset-sink/src/main/java/org/apache/flume/sink/kite/policy/FailurePolicy.java:78
↓ 3 callers
Method
sync
()
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/HDFSWriter.java:41
↓ 3 callers
Method
testDocumentTypesInternal
(String... files)
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/test/java/org/apache/flume/sink/solr/morphline/TestMorphlineSolrSink.java:328
↓ 3 callers
Method
unregisterCurrentStream
()
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/AbstractHDFSWriter.java:121
↓ 3 callers
Method
validateMiniParse
(EventDeserializer des)
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/test/java/org/apache/flume/sink/solr/morphline/TestBlobDeserializer.java:103
↓ 3 callers
Method
write
Parse the event using the entity parser and write the entity to the dataset. @param event The event to write @throws EventDeliveryException An error
code/flume-ng-sinks/flume-dataset-sink/src/main/java/org/apache/flume/sink/kite/DatasetSink.java:363
↓ 2 callers
Method
addSimpleField
(XContentBuilder builder, String fieldName, byte[] data)
code/flume-ng-sinks/flume-ng-elasticsearch-sink/src/main/java/org/apache/flume/sink/elasticsearch/ContentBuilderUtil.java:51
↓ 2 callers
Method
afterReopen
()
code/flume-ng-sinks/flume-hdfs-sink/src/main/java/org/apache/flume/sink/hdfs/AvroEventSerializer.java:105
↓ 2 callers
Method
apply
(@Nullable GenericRecord rec)
code/flume-ng-sinks/flume-dataset-sink/src/test/java/org/apache/flume/sink/kite/TestDatasetSink.java:164
↓ 2 callers
Method
assertEqualsEventList
(List<Event> x, List<Event> y)
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/test/java/org/apache/flume/sink/solr/morphline/TestMorphlineInterceptor.java:160
↓ 2 callers
Method
borrowFromPool
()
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/main/java/org/apache/flume/sink/solr/morphline/MorphlineInterceptor.java:91
↓ 2 callers
Method
call
()
code/flume-ng-sinks/flume-hive-sink/src/main/java/org/apache/flume/sink/hive/HiveWriter.java:475
↓ 2 callers
Method
checkIfChannelExceptionAndThrow
(Throwable e)
code/flume-ng-sinks/flume-ng-hbase-sink/src/main/java/org/apache/flume/sink/hbase/AsyncHBaseSink.java:666
↓ 2 callers
Method
closeTxnBatch
()
code/flume-ng-sinks/flume-hive-sink/src/main/java/org/apache/flume/sink/hive/HiveWriter.java:405
↓ 2 callers
Method
compare
(SearchHit o1, SearchHit o2)
code/flume-ng-sinks/flume-ng-elasticsearch-sink/src/test/java/org/apache/flume/sink/elasticsearch/AbstractElasticSearchSinkTest.java:146
↓ 2 callers
Method
configure
(Context context)
code/flume-ng-sinks/flume-irc-sink/src/main/java/org/apache/flume/sink/irc/IRCSink.java:127
↓ 2 callers
Method
configure
(Context arg0)
code/flume-ng-sinks/flume-ng-elasticsearch-sink/src/test/java/org/apache/flume/sink/elasticsearch/TestElasticSearchIndexRequestBuilderFactory.java:204
↓ 2 callers
Method
configure
(Context context)
code/flume-ng-sinks/flume-ng-elasticsearch-sink/src/main/java/org/apache/flume/sink/elasticsearch/ElasticSearchLogStashEventSerializer.java:136
↓ 2 callers
Method
configure
(Context context)
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/main/java/org/apache/flume/sink/solr/morphline/UUIDInterceptor.java:108
↓ 2 callers
Method
configureHostnames
(String[] hostNames)
code/flume-ng-sinks/flume-ng-elasticsearch-sink/src/main/java/org/apache/flume/sink/elasticsearch/client/ElasticSearchTransportClient.java:133
↓ 2 callers
Method
createConnection
()
code/flume-ng-sinks/flume-irc-sink/src/main/java/org/apache/flume/sink/irc/IRCSink.java:153
↓ 2 callers
Method
createTopic
(String topicName, int numPartitions)
code/flume-ng-sinks/flume-ng-kafka-sink/src/test/java/org/apache/flume/sink/kafka/TestKafkaSink.java:519
↓ 2 callers
Method
deleteTopic
(String topicName)
code/flume-ng-sinks/flume-ng-kafka-sink/src/test/java/org/apache/flume/sink/kafka/TestKafkaSink.java:529
↓ 2 callers
Method
dirCleanup
()
code/flume-ng-sinks/flume-hdfs-sink/src/test/java/org/apache/flume/sink/hdfs/TestHDFSEventSink.java:88
↓ 2 callers
Method
doTestMultipleBatchesBatchIncrements
(boolean coalesce)
code/flume-ng-sinks/flume-ng-hbase-sink/src/test/java/org/apache/flume/sink/hbase/TestAsyncHBaseSink.java:289
↓ 2 callers
Method
doTestTextBatchAppend
(boolean useRawLocalFileSystem)
code/flume-ng-sinks/flume-hdfs-sink/src/test/java/org/apache/flume/sink/hdfs/TestHDFSEventSink.java:135
↓ 2 callers
Method
findUnusedTopic
()
code/flume-ng-sinks/flume-ng-kafka-sink/src/test/java/org/apache/flume/sink/kafka/TestKafkaSink.java:537
↓ 2 callers
Method
generateEvents
Add number of Events corresponding to counts to the events list. @param events Destination list. @param counts How many events to generate for each ro
code/flume-ng-sinks/flume-ng-hbase-sink/src/test/java/org/apache/flume/sink/hbase/TestHBaseSink.java:714
↓ 2 callers
Method
getConfig
()
code/flume-ng-sinks/flume-ng-hbase-sink/src/main/java/org/apache/flume/sink/hbase/HBaseSink.java:300
↓ 2 callers
Method
getContentBuilder
(Event event)
code/flume-ng-sinks/flume-ng-elasticsearch-sink/src/main/java/org/apache/flume/sink/elasticsearch/ElasticSearchLogStashEventSerializer.java:76
↓ 2 callers
Method
getContentType
()
code/flume-ng-sinks/flume-ng-morphline-solr-sink/src/test/java/org/apache/flume/sink/solr/morphline/FlumeHttpServletRequestWrapper.java:202
↓ 2 callers
Method
getEventSerializer
()
code/flume-ng-sinks/flume-ng-elasticsearch-sink/src/main/java/org/apache/flume/sink/elasticsearch/ElasticSearchSink.java:164
← previous
next →
101–200 of 982, ranked by callers