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
Method
getOffsetsToCommit
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:241
Method
getParallelInstanceId
Get the number of the subtask that RocketMQSink is running on. The numbering starts from 0 and goes up to parallelism-1. (parallelism as returned by {
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContext.java:38
Method
getParallelInstanceId
()
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContextImpl.java:43
Method
getProducedType
()
src/test/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceTest.java:57
Method
getProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:654
Method
getProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleStringDeserializationSchema.java:35
Method
getProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleTupleDeserializationSchema.java:36
Method
getProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:120
Method
getProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/ForwardMessageExtDeserialization.java:33
Method
getProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleKeyValueDeserializationSchema.java:63
Method
getProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:189
Method
getProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:395
Method
getProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQSchemaWrapper.java:30
Method
getProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchemaWrapper.java:53
Method
getProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQRowDeserializationSchema.java:81
Method
getProducerGroup
Gets the consumer group of the consumer. @return the consumer group of the consumer
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducer.java:43
Method
getProducerGroup
()
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:154
Method
getProperties
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:138
Method
getProperties
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/BytesMessage.java:39
Method
getProperties
Get the option value by a prefix. We would return an empty map if the option doesn't exist.
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfiguration.java:55
Method
getQueueId
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:73
Method
getQueueId
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:93
Method
getQueueId
()
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:49
Method
getQueueOffset
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:81
Method
getQueueOffset
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:98
Method
getReSendInitAssign
()
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceDetectEvent.java:27
Method
getScanRuntimeProvider
(ScanContext scanContext)
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:119
Method
getSinkRuntimeProvider
( DynamicTableSink.Context context)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:157
Method
getSplitSerializer
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:179
Method
getStoreSize
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:113
Method
getStrategyName
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/ConsistentHashAllocateStrategy.java:29
Method
getStrategyName
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AverageAllocateStrategy.java:12
Method
getStrategyName
Allocate strategy name @return Current strategy name
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AllocateStrategy.java:36
Method
getStrategyName
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/BroadcastAllocateStrategy.java:29
Method
getTag
(T tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/TopicSelector.java:25
Method
getTag
Get the tag of the message, which is used for filtering. @return the message tag
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:66
Method
getTag
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:103
Method
getTopic
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/selector/SimpleTopicSelectorTest.java:29
Method
getTopic
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/selector/DefaultTopicSelectorTest.java:26
Method
getTopic
(T tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/TopicSelector.java:23
Method
getTopic
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:57
Method
getTopic
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:83
Method
getTopic
()
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:57
Method
getUserCodeClassLoader
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:118
Method
getValue
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:100
Method
getVersion
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittableSerializer.java:33
Method
handleSourceEvent
(int taskId, SourceEvent sourceEvent)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:251
Method
handleSourceEvents
(SourceEvent sourceEvent)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:166
Method
handleSourceQueueChange
(Set<MessageQueue> latestSet, Throwable t)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:300
Method
handleSplitChanges
Mark partition splits initialized by {@link RocketMQSourceEnumerator#initializeSourceSplits(SourceChangeResult)} as pending and try to assign pending
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:385
Method
handleSplitRequest
(int subtaskId, @Nullable String requesterHostname)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:196
Method
initializeState
(FunctionInitializationContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:237
Method
initializeState
called every time the user-defined function is initialized, be that when the function is first initialized or be that when the function is actually re
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:623
Method
initializedState
(RocketMQSourceSplit partitionSplit)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:148
Method
invoke
(Message input, Context context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:109
Method
invoke
(RowData rowData, Context context)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataSink.java:49
Method
isEnableSchemaEvolution
RocketMQ can check the schema and upgrade the schema automatically. If you enable this option, we wouldn't serialize the record into bytes, we send an
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContext.java:47
Method
isEnableSchemaEvolution
()
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContextImpl.java:53
Method
isRunning
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceTest.java:109
Method
latest
Get an {@link OffsetsSelector} which initializes the offsets to the latest offsets of each partition. @return an {@link OffsetsSelector} which initia
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelector.java:138
Method
listReadableMetadata
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:150
Method
listWritableMetadata
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:163
Method
main
(String[] args)
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleConsumer.java:37
Method
main
(String[] args)
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleProducer.java:40
Method
main
(String[] args)
src/test/java/org/apache/flink/connector/rocketmq/example/LegacyConnectorExample.java:81
Method
main
(String[] args)
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleAdmin.java:100
Method
main
(String[] args)
src/test/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkTest.java:50
Method
main
(String[] args)
src/test/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceTest.java:31
Method
markActive
()
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:86
Method
markIdle
()
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:82
Method
maxOffsets
(Collection<MessageQueue> messageQueues)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:417
Method
minOffsets
(Collection<MessageQueue> messageQueues)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:388
Method
nextRecordFromSplit
Gets the next record from the current split. Returns null if no more records are left in this split.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:328
Method
nextSplit
Moves to the next split. This method is also called initially to move to the first split. Returns null, if no splits are left.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:309
Method
notifyAssignResult
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:641
Method
notifyCheckpointComplete
(long checkpointId)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:659
Method
notifyCheckpointComplete
(long checkpointId)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:130
Method
offsets
Get an {@link OffsetsSelector} which initializes the offsets to the specified offsets. @param offsets the specified offsets for each partition. @retu
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelector.java:150
Method
offsetsForTimes
( Map<MessageQueue, Long> messageQueueWithTimeMap)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:446
Method
onException
(Throwable throwable)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:137
Method
onSplitFinished
(Map<String, RocketMQSourceSplitState> finishedSplitIds)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:93
Method
onSuccess
(SendResult sendResult)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:128
Method
open
(Configuration parameters)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:77
Method
open
(Configuration parameters)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:147
Method
open
(Configuration configuration)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataSink.java:38
Method
open
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataConverter.java:102
Method
open
(InitializationContext context)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:123
Method
open
(DeserializationSchema.InitializationContext context)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchemaWrapper.java:45
Method
open
(InitializationContext context)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQRowDeserializationSchema.java:63
Method
open
()
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:108
Method
optionalOptions
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSinkFactory.java:73
Method
optionalOptions
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQDynamicTableSourceFactory.java:65
Method
override
Override the option with the given value. It will not check the existed option comparing to {@link #set(ConfigOption, Object)}.
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigBuilder.java:130
Method
pause
(Collection<MessageQueue> messageQueues)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:223
Method
poll
(Duration timeout)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:189
Method
prepareCommit
()
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/RocketMQWriter.java:113
Method
processElement
(Tuple2<String, String> tuple, Context ctx, Collector<Message> out)
src/main/java/org/apache/flink/connector/rocketmq/legacy/function/SinkMapFunction.java:40
Method
processElement
( Tuple2<String, String> value, Context ctx, Collector<Tuple2<String, String>> out)
src/main/java/org/apache/flink/connector/rocketmq/legacy/function/SourceMapFunction.java:27
Method
processTime
Returns the current process time in flink.
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContext.java:51
Method
processTime
()
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContextImpl.java:58
← previous
next →
501–600 of 733, ranked by callers