Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/apache/rocketmq-flink
/ functions
Functions
733 in github.com/apache/rocketmq-flink
⨍
Functions
733
◇
Types & classes
161
↓ 1 callers
Method
setInitFinish
(boolean initFinish)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:258
↓ 1 callers
Method
setLineDelimiter
(String lineDelimiter)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:500
↓ 1 callers
Method
setMetadataConverters
(MetadataConverter[] metadataConverters)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:515
↓ 1 callers
Method
setOffsetMsgId
(String offsetMsgId)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:101
↓ 1 callers
Method
setProperties
Set arbitrary properties for the RocketMQSink and RocketMQ Consumer. The valid keys can be found in {@link RocketMQSinkOptions} and {@link RocketMQOpt
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:136
↓ 1 callers
Method
setQueueId
(Integer queueId)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:77
↓ 1 callers
Method
setQueueId
(int queueId)
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:37
↓ 1 callers
Method
setQueueOffset
(Long queueOffset)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:85
↓ 1 callers
Method
setSerializer
Sets the {@link RocketMQSerializationSchema} that transforms incoming records to {@link org.apache.rocketmq.common.message.MessageExt}s. @param seria
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:148
↓ 1 callers
Method
setSourceOutput
(SourceOutput<T> sourceOutput)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitter.java:69
↓ 1 callers
Method
setStartFromEarliest
consume from the min offset at every restart with no state
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:442
↓ 1 callers
Method
setStartFromSpecificOffsets
consume from the specific offset. Group offsets is enable while the broker didn't specify offset.
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:481
↓ 1 callers
Method
setStartFromTimeStamp
consume from the closest offset
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:454
↓ 1 callers
Method
setTimestamp
(long timestamp)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitter.java:73
↓ 1 callers
Method
setTopic
(String topic)
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:41
↓ 1 callers
Method
setTopics
Set a list of topics the RocketMQSource should consume from. All the topics in the list should have existed in the RocketMQ cluster. Otherwise, an exc
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:93
↓ 1 callers
Method
setUnbounded
(OffsetsSelector stoppingOffsetsSelector)
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:116
↓ 1 callers
Method
splitId
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:142
↓ 1 callers
Method
toBigDecimal
(byte[] bytes)
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteUtils.java:205
↓ 1 callers
Method
toBoolean
Convert a byte array to a boolean. @param b array @return True or false.
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteUtils.java:64
↓ 1 callers
Method
toDouble
Parse a byte array to double. @param bytes byte array @return Return double made from passed bytes.
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteUtils.java:140
↓ 1 callers
Method
toFloat
Presumes float encoded as IEEE 754 floating-point "single format". @param bytes byte array @return Float made from passed byte array.
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteUtils.java:119
↓ 1 callers
Method
toShort
Converts a byte array to a short value. @param bytes byte array @return the short value
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteUtils.java:161
↓ 1 callers
Method
toSplitId
(MessageQueue messageQueue)
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:124
↓ 1 callers
Method
unregisterMessageQueue
(MessageQueue messageQueue)
src/main/java/org/apache/flink/connector/rocketmq/source/metrics/RocketMQSourceReaderMetrics.java:49
↓ 1 callers
Method
updateMessageQueueOffset
(MessageQueue mq, long offset)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:488
↓ 1 callers
Method
withAsync
(boolean async)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:183
↓ 1 callers
Method
withBatchFlushOnCheckpoint
(boolean batchFlushOnCheckpoint)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:188
↓ 1 callers
Method
withBatchSize
(int batchSize)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:193
Method
AllocateStrategyFactory
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AllocateStrategyFactory.java:32
Method
BoundedOutOfOrdernessGenerator
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/BoundedOutOfOrdernessGenerator.java:31
Method
BoundedOutOfOrdernessGeneratorPerQueue
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/BoundedOutOfOrdernessGeneratorPerQueue.java:34
Method
Builder
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:297
Method
Builder
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:468
Method
DefaultTopicSelector
(final String topicName, final String tagName)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/DefaultTopicSelector.java:23
Method
InnerConsumerImpl
(Configuration configuration)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:68
Method
InnerProducerImpl
(Configuration configuration)
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:65
Method
MessageViewExt
(MessageExt messageExt)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:59
Method
MetadataCollector
(boolean hasMetadata, MetadataConverter[] metadataConverters)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:420
Method
OffsetVerification
(InlineElement desc)
src/main/java/org/apache/flink/connector/rocketmq/source/config/OffsetVerification.java:40
Method
OffsetsSelectorBySpecified
( Map<MessageQueue, Long> initialOffsets, OffsetResetStrategy offsetResetStrategy)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorBySpecified.java:37
Method
OffsetsSelectorByStrategy
( ConsumeFromWhere consumeFromWhere, OffsetResetStrategy offsetResetStrategy)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorByStrategy.java:34
Method
OffsetsSelectorByTimestamp
(long startingTimestamp)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorByTimestamp.java:32
Method
ReadableMetadata
(String key, DataType dataType, MetadataConverter converter)
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:257
Method
RemotingOffsetsRetrieverImpl
(InnerConsumer innerConsumer)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:350
Method
RetryUtil
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RetryUtil.java:36
Method
RocketMQCatalog
( String catalogName, String database, String namesrvAddr, String schemaRegistryUrl)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:89
Method
RocketMQCatalogFactoryOptions
()
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactoryOptions.java:51
Method
RocketMQCommitter
(Configuration configuration)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/RocketMQCommitter.java:47
Method
RocketMQConfigValidator
( List<Set<ConfigOption<?>>> conflictOptions, Set<ConfigOption<?>> requiredOptions)
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigValidator.java:51
Method
RocketMQConfiguration
Creates a new RocketMQConfiguration, which holds a copy of the given configuration that can't be altered. @param config The configuration with the or
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfiguration.java:48
Method
RocketMQDeserializationSchemaWrapper
(DeserializationSchema<T> deserializationSchema)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchemaWrapper.java:41
Method
RocketMQDynamicTableSink
( DescriptorProperties properties, TableSchema schema, String topic,
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:71
Method
RocketMQRecordEmitter
(RocketMQDeserializationSchema<T> deserializationSchema)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitter.java:36
Method
RocketMQRecordsWithSplitIds
(RocketMQSourceReaderMetrics readerMetrics)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:288
Method
RocketMQRowDataConverter
( String topic, String tag, String dynamicColumn, String field
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataConverter.java:73
Method
RocketMQRowDataSink
(RocketMQSink sink, RocketMQRowDataConverter converter)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataSink.java:33
Method
RocketMQRowDeserializationSchema
( TableSchema tableSchema, Map<String, String> properties, boolean hasMeta
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQRowDeserializationSchema.java:49
Method
RocketMQScanTableSource
( long pollTime, DescriptorProperties properties, TableSchema schema,
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:76
Method
RocketMQSink
(Properties props)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:73
Method
RocketMQSink
( Configuration configuration, MessageQueueSelector messageQueueSelector,
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSink.java:39
Method
RocketMQSinkBuilder
()
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:58
Method
RocketMQSinkContextImpl
(InitContext initContext, Configuration configuration)
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContextImpl.java:35
Method
RocketMQSource
( OffsetsSelector startingOffsetsSelector, OffsetsSelector stoppingOffsetsSelector,
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:77
Method
RocketMQSourceBuilder
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:58
Method
RocketMQSourceEnumState
(Set<MessageQueue> currentSplitAssignment)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumState.java:33
Method
RocketMQSourceEnumerator
( OffsetsSelector startingOffsetsSelector, OffsetsSelector stoppingOffsetsSelector,
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:110
Method
RocketMQSourceFetcherManager
Creates a new SplitFetcherManager with a single I/O threads. @param elementsQueue The queue that is used to hand over data from the I/O thread (the
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceFetcherManager.java:55
Method
RocketMQSourceFunction
(KeyValueDeserializationSchema<OUT> schema, Properties props)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:142
Method
RocketMQSourceReader
( FutureCompletingBlockingQueue<RecordsWithSplitIds<MessageView>> elementsQueue, Rocke
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:75
Method
RocketMQSourceReaderMetrics
(SourceReaderMetricGroup sourceReaderMetricGroup)
src/main/java/org/apache/flink/connector/rocketmq/source/metrics/RocketMQSourceReaderMetrics.java:45
Method
RocketMQSourceSplit
( MessageQueue messageQueue, long startingOffset, long stoppingOffset)
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:43
Method
RocketMQSourceSplitState
(RocketMQSourceSplit partitionSplit)
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplitState.java:26
Method
RocketMQSplitReader
( Configuration configuration, SourceReaderContext sourceReaderContext, Ro
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:93
Method
RocketMQWriter
( Configuration configuration, MessageQueueSelector messageQueueSelector,
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/RocketMQWriter.java:66
Method
RowDeserializationSchema
( TableSchema tableSchema, DirtyDataStrategy formatErrorStrategy, DirtyDat
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:79
Method
RowKeyValueDeserializationSchema
( TableSchema tableSchema, DirtyDataStrategy formatErrorStrategy, DirtyDat
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:70
Method
SendCommittable
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:45
Method
SimpleAdmin
()
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleAdmin.java:41
Method
SimpleKeyValueDeserializationSchema
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleKeyValueDeserializationSchema.java:34
Method
SimpleKeyValueSerializationSchema
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleKeyValueSerializationSchema.java:29
Method
SimpleTopicSelector
SimpleTopicSelector Constructor. @param topicFieldName field name used for selecting topic @param defaultTopicName default field name used for select
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/SimpleTopicSelector.java:42
Method
SinkMapFunction
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/function/SinkMapFunction.java:33
Method
SourceChangeResult
(Set<MessageQueue> increaseSet, Set<MessageQueue> decreaseSet)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:674
Method
SourceSplitChangeResult
(Set<RocketMQSourceSplit> increaseSet)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:698
Method
TestingEmitterOutput
()
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:70
Method
TimeLagWatermarkGenerator
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/TimeLagWatermarkGenerator.java:32
Method
WaterMarkForAll
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/WaterMarkForAll.java:28
Method
WaterMarkPerQueue
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/WaterMarkPerQueue.java:34
Method
WritableMetadata
( String key, DataType dataType, RocketMQRowDataConverter.Meta
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:300
Method
addReader
(int subtaskId)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:219
Method
addSplitsBack
Add a split back to the split enumerator. It will only happen when a {@link SourceReader} fails and there are splits assigned to it after the last suc
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:208
Method
applyReadableMetadata
(List<String> metadataKeys, DataType producedDataType)
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:158
Method
applyWritableMetadata
(List<String> metadataKeys, DataType consumedDataType)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:171
Method
asSummaryString
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:188
Method
assign
(Collection<MessageQueue> messageQueues)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:175
Method
assignment
()
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:180
Method
averageAllocateStrategyTest
()
src/test/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AverageAllocateStrategyTest.java:22
Method
averagesAllocateStrategyTest
()
src/test/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AverageAllocateStrategyTest.java:43
Method
broadcastAllocateStrategyTest
()
src/test/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/BroadcastAllocateStrategyTest.java:22
← previous
next →
301–400 of 733, ranked by callers