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
↓ 79 callers
Method
get
Get an option value from the given config, convert it into a new value instance.
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfiguration.java:80
↓ 38 callers
Method
setProperty
(String key, String value)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/BytesMessage.java:51
↓ 31 callers
Method
equals
(Object obj)
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:159
↓ 27 callers
Method
getTopic
Get the topic that the message belongs to. @return the topic
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:38
↓ 26 callers
Method
isEmpty
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:687
↓ 26 callers
Method
set
Add a config option with a not null value. The config key shouldn't be duplicated. @param option Config option instance, contains key & type definiti
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigBuilder.java:66
↓ 25 callers
Method
getQueueId
Get the ID of the queue that the message is stored in. @return the queue ID
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:52
↓ 23 callers
Method
getProperty
(String key)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/BytesMessage.java:47
↓ 21 callers
Method
getMessageQueue
(RocketMQSourceSplit split)
src/main/java/org/apache/flink/connector/rocketmq/source/util/UtilAll.java:38
↓ 20 callers
Method
getBrokerName
Get the name of the broker that handles the message. @return the broker name
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:45
↓ 15 callers
Method
contains
Validate if the config has a existed option.
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigBuilder.java:46
↓ 15 callers
Method
getStartingOffset
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:108
↓ 15 callers
Method
getStoppingOffset
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:112
↓ 12 callers
Method
collect
(T record)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitter.java:61
↓ 12 callers
Method
toString
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:143
↓ 11 callers
Method
getMetricGroup
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:113
↓ 11 callers
Method
getValue
(BytesMessage message, String[] data, String line, int index)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:252
↓ 11 callers
Method
setFieldValue
(Object obj, String fieldName, Object value)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/TestUtils.java:22
↓ 10 callers
Method
getInteger
(Properties props, String key, int defaultValue)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RocketMQUtils.java:31
↓ 10 callers
Method
getIsIncrease
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:120
↓ 10 callers
Method
getMsgId
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:89
↓ 10 callers
Method
getQueueDescription
(MessageQueue mq)
src/main/java/org/apache/flink/connector/rocketmq/source/util/UtilAll.java:32
↓ 9 callers
Method
getBody
Get the body of the message. @return the message body
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:87
↓ 9 callers
Method
getFieldValue
(Object obj, String fieldName)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/TestUtils.java:32
↓ 9 callers
Method
send
Sends the message to the messaging system and returns a Future for the send operation. @param message the message to be sent @return a Future for the
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducer.java:53
↓ 9 callers
Method
start
Starts the inner consumer.
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducer.java:36
↓ 8 callers
Method
build
Build the {@link RocketMQSource}. @return a RocketMQSource with the settings made for this builder.
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:159
↓ 8 callers
Method
getBoolean
(Properties props, String key, boolean defaultValue)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RocketMQUtils.java:39
↓ 8 callers
Method
getLong
(Properties props, String key, long defaultValue)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RocketMQUtils.java:35
↓ 8 callers
Method
getTransactionId
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:105
↓ 7 callers
Method
getQueueOffset
Get the offset of the message within the queue. @return the queue offset
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:59
↓ 7 callers
Method
getValue
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:87
↓ 6 callers
Method
getBrokerName
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:100
↓ 6 callers
Method
getQueueId
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:104
↓ 6 callers
Method
getTopic
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:96
↓ 6 callers
Method
lock
()
src/main/java/org/apache/flink/connector/rocketmq/common/lock/SpinLock.java:25
↓ 6 callers
Method
setConfig
Set an arbitrary property for the RocketMQ source. The valid keys can be found in {@link RocketMQSourceOptions}. <p>Make sure the option could be set
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:112
↓ 6 callers
Method
toString
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/WaterMarkPerQueue.java:55
↓ 6 callers
Method
toString
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:129
↓ 6 callers
Method
unlock
()
src/main/java/org/apache/flink/connector/rocketmq/common/lock/SpinLock.java:32
↓ 5 callers
Method
assignment
Returns a set of message queues that are assigned to the current consumer. The assignment is typically performed by a message broker and may change dy
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:61
↓ 5 callers
Method
getInstanceName
(String... args)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RocketMQUtils.java:48
↓ 5 callers
Method
getKeys
Get the keys of the message, which are used for partitioning and indexing. @return the message keys
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:73
↓ 5 callers
Method
getMessageQueueOffsets
This method retrieves the current offsets for a collection of {@link MessageQueue}s using a provided {@link MessageQueueOffsetsRetriever}. @param mes
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelector.java:48
↓ 5 callers
Method
isRunning
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/RunningChecker.java:24
↓ 5 callers
Method
toLong
Converts a byte array to a long value. @param bytes array @return the long value
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteUtils.java:88
↓ 5 callers
Method
validate
Validate offsets initializer with properties of RocketMQ source. @param properties Properties of RocketMQ source @throws IllegalStateException if val
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsValidator.java:16
↓ 4 callers
Method
allocate
( Collection<RocketMQSourceSplit> mqAll, int parallelism)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AverageAllocateStrategy.java:17
↓ 4 callers
Method
createTableSource
( Map<String, String> options, Configuration conf)
src/test/java/org/apache/flink/connector/rocketmq/source/table/RocketMQDynamicTableSourceFactoryTest.java:111
↓ 4 callers
Method
deserialize
(int version, byte[] serialized)
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQPartitionSplitSerializer.java:56
↓ 4 callers
Method
deserialize
( String value, ByteSerializer.ValueType type, DataType dataType,
src/main/java/org/apache/flink/connector/rocketmq/source/util/StringSerializer.java:41
↓ 4 callers
Method
explainWrongLengthOrOffset
( final byte[] bytes, final int offset, final int length, final int expectedLength)
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteUtils.java:187
↓ 4 callers
Method
fetchMessageQueues
Fetch message queues of the topic. @param topic topic list @return key is topic, values are message queue collections
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:44
↓ 4 callers
Method
get
Get an option-related config value. We would return the default config value defined in the option if no value existed instead. @param key Config opt
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigBuilder.java:56
↓ 4 callers
Method
getCurrentOffset
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplitState.java:36
↓ 4 callers
Method
getCurrentSplitAssignment
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumState.java:37
↓ 4 callers
Method
getMessageQueue
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:113
↓ 4 callers
Method
getProperties
Get the properties of the message, which are set by the producer. @return the message properties
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:116
↓ 4 callers
Method
getSplits
()
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceInitAssignEvent.java:33
↓ 4 callers
Method
getTypeIndex
(Class<?> clazz)
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteSerializer.java:86
↓ 4 callers
Method
poll
Fetch data for the topics or partitions specified using assign API @return list of message, can be null.
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:68
↓ 4 callers
Method
report
(long timeDelta, long batchSize)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:81
↓ 4 callers
Method
serialize
(RocketMQSourceSplit split)
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQPartitionSplitSerializer.java:41
↓ 4 callers
Method
setFieldIncrementStrategy
(DirtyDataStrategy fieldIncrementStrategy)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:314
↓ 4 callers
Method
setFieldIncrementStrategy
(DirtyDataStrategy fieldIncrementStrategy)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:485
↓ 4 callers
Method
setFieldMissingStrategy
(DirtyDataStrategy fieldMissingStrategy)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:309
↓ 4 callers
Method
setFieldMissingStrategy
(DirtyDataStrategy fieldMissingStrategy)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:480
↓ 4 callers
Method
setFormatErrorStrategy
(DirtyDataStrategy formatErrorStrategy)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:304
↓ 4 callers
Method
setFormatErrorStrategy
(DirtyDataStrategy formatErrorStrategy)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:475
↓ 4 callers
Method
setRunning
(boolean running)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RunningChecker.java:28
↓ 4 callers
Method
start
start inner consumer
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:33
↓ 3 callers
Method
asSummaryString
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:201
↓ 3 callers
Method
call
(Callable<T> callable, String errorMsg)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RetryUtil.java:46
↓ 3 callers
Method
collect
(RowData physicalRow)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:425
↓ 3 callers
Method
committedOffsets
Get an {@link OffsetsSelector} which initializes the offsets to the committed offsets. An exception will be thrown at runtime if there is no committed
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelector.java:89
↓ 3 callers
Method
createDynamicTableSink
(Map<String, String> options)
src/test/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSinkFactoryTest.java:97
↓ 3 callers
Method
flush
(boolean endOfInput)
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/RocketMQWriter.java:108
↓ 3 callers
Method
flushSync
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:221
↓ 3 callers
Method
getAclRpcHook
()
src/test/java/org/apache/flink/connector/rocketmq/example/ConnectorConfig.java:50
↓ 3 callers
Method
getCheckpoint
()
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:45
↓ 3 callers
Method
getData
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/BytesMessage.java:31
↓ 3 callers
Method
getHeaderValue
(BytesMessage message, int index)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:247
↓ 3 callers
Method
getIncreaseSet
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:679
↓ 3 callers
Method
hashCode
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:154
↓ 3 callers
Method
maxOffsets
List max offsets for the specified MessageQueues.
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelector.java:75
↓ 3 callers
Method
open
Initialization method for the schema. It is called before the actual working methods {@link #deserialize} and thus suitable for one time setup work.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/QueryableSchema.java:41
↓ 3 callers
Method
seek
Overrides the fetch offsets that the consumer will use on the next poll. If this method is invoked for the same message queue more than once, the late
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:95
↓ 3 callers
Method
sendSplitChangesToRemote
(Set<Integer> pendingReaders)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:541
↓ 3 callers
Method
setProperties
(Map<String, String> props)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/BytesMessage.java:43
↓ 3 callers
Method
setTableSchema
(TableSchema tableSchema)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:470
↓ 3 callers
Method
setTopic
(String topic)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:61
↓ 3 callers
Method
toInt
Converts a byte array to an int value. @param bytes byte array @return the int value
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteUtils.java:33
↓ 2 callers
Method
addFinishedSplit
(String splitId)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:297
↓ 2 callers
Method
allocate
Average Hashing queue algorithm Refer: org.apache.rocketmq.client.consumer.rebalance.AllocateStrategyByAveragely
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RocketMQUtils.java:59
↓ 2 callers
Method
allocate
Allocates RocketMQ source splits to Flink tasks based on the selected allocation strategy. @param mqAll a collection of all available RocketMQ source
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AllocateStrategy.java:45
↓ 2 callers
Method
assign
Manually assign a list of message queues to this consumer. This interface does not allow for incremental assignment and will replace the previous assi
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:52
↓ 2 callers
Method
build
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:573
↓ 2 callers
Method
buildAclRPCHook
Build credentials for client. @param props @return
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQConfig.java:183
↓ 2 callers
Method
buildCommonConfigs
Build Common Configs. @param props Properties @param client ClientConfig
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQConfig.java:166
↓ 2 callers
Method
builder
Return the builder for building {@link RocketMQConfigValidator}.
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigValidator.java:78
next →
1–100 of 733, ranked by callers