MCPcopy Create free account

hub / github.com/apache/rocketmq-flink / functions

Functions733 in github.com/apache/rocketmq-flink

↓ 79 callersMethodget
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 callersMethodsetProperty
(String key, String value)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/BytesMessage.java:51
↓ 31 callersMethodequals
(Object obj)
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:159
↓ 27 callersMethodgetTopic
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 callersMethodisEmpty
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:687
↓ 26 callersMethodset
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 callersMethodgetQueueId
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 callersMethodgetProperty
(String key)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/BytesMessage.java:47
↓ 21 callersMethodgetMessageQueue
(RocketMQSourceSplit split)
src/main/java/org/apache/flink/connector/rocketmq/source/util/UtilAll.java:38
↓ 20 callersMethodgetBrokerName
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 callersMethodcontains
Validate if the config has a existed option.
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigBuilder.java:46
↓ 15 callersMethodgetStartingOffset
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:108
↓ 15 callersMethodgetStoppingOffset
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:112
↓ 12 callersMethodcollect
(T record)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitter.java:61
↓ 12 callersMethodtoString
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:143
↓ 11 callersMethodgetMetricGroup
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:113
↓ 11 callersMethodgetValue
(BytesMessage message, String[] data, String line, int index)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:252
↓ 11 callersMethodsetFieldValue
(Object obj, String fieldName, Object value)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/TestUtils.java:22
↓ 10 callersMethodgetInteger
(Properties props, String key, int defaultValue)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RocketMQUtils.java:31
↓ 10 callersMethodgetIsIncrease
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:120
↓ 10 callersMethodgetMsgId
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:89
↓ 10 callersMethodgetQueueDescription
(MessageQueue mq)
src/main/java/org/apache/flink/connector/rocketmq/source/util/UtilAll.java:32
↓ 9 callersMethodgetBody
Get the body of the message. @return the message body
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:87
↓ 9 callersMethodgetFieldValue
(Object obj, String fieldName)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/TestUtils.java:32
↓ 9 callersMethodsend
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 callersMethodstart
Starts the inner consumer.
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducer.java:36
↓ 8 callersMethodbuild
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 callersMethodgetBoolean
(Properties props, String key, boolean defaultValue)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RocketMQUtils.java:39
↓ 8 callersMethodgetLong
(Properties props, String key, long defaultValue)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RocketMQUtils.java:35
↓ 8 callersMethodgetTransactionId
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:105
↓ 7 callersMethodgetQueueOffset
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 callersMethodgetValue
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:87
↓ 6 callersMethodgetBrokerName
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:100
↓ 6 callersMethodgetQueueId
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:104
↓ 6 callersMethodgetTopic
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:96
↓ 6 callersMethodlock
()
src/main/java/org/apache/flink/connector/rocketmq/common/lock/SpinLock.java:25
↓ 6 callersMethodsetConfig
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 callersMethodtoString
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/WaterMarkPerQueue.java:55
↓ 6 callersMethodtoString
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:129
↓ 6 callersMethodunlock
()
src/main/java/org/apache/flink/connector/rocketmq/common/lock/SpinLock.java:32
↓ 5 callersMethodassignment
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 callersMethodgetInstanceName
(String... args)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RocketMQUtils.java:48
↓ 5 callersMethodgetKeys
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 callersMethodgetMessageQueueOffsets
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 callersMethodisRunning
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/RunningChecker.java:24
↓ 5 callersMethodtoLong
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 callersMethodvalidate
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 callersMethodallocate
( Collection<RocketMQSourceSplit> mqAll, int parallelism)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AverageAllocateStrategy.java:17
↓ 4 callersMethodcreateTableSource
( Map<String, String> options, Configuration conf)
src/test/java/org/apache/flink/connector/rocketmq/source/table/RocketMQDynamicTableSourceFactoryTest.java:111
↓ 4 callersMethoddeserialize
(int version, byte[] serialized)
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQPartitionSplitSerializer.java:56
↓ 4 callersMethoddeserialize
( String value, ByteSerializer.ValueType type, DataType dataType,
src/main/java/org/apache/flink/connector/rocketmq/source/util/StringSerializer.java:41
↓ 4 callersMethodexplainWrongLengthOrOffset
( 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 callersMethodfetchMessageQueues
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 callersMethodget
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 callersMethodgetCurrentOffset
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplitState.java:36
↓ 4 callersMethodgetCurrentSplitAssignment
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumState.java:37
↓ 4 callersMethodgetMessageQueue
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:113
↓ 4 callersMethodgetProperties
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 callersMethodgetSplits
()
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceInitAssignEvent.java:33
↓ 4 callersMethodgetTypeIndex
(Class<?> clazz)
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteSerializer.java:86
↓ 4 callersMethodpoll
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 callersMethodreport
(long timeDelta, long batchSize)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:81
↓ 4 callersMethodserialize
(RocketMQSourceSplit split)
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQPartitionSplitSerializer.java:41
↓ 4 callersMethodsetFieldIncrementStrategy
(DirtyDataStrategy fieldIncrementStrategy)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:314
↓ 4 callersMethodsetFieldIncrementStrategy
(DirtyDataStrategy fieldIncrementStrategy)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:485
↓ 4 callersMethodsetFieldMissingStrategy
(DirtyDataStrategy fieldMissingStrategy)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:309
↓ 4 callersMethodsetFieldMissingStrategy
(DirtyDataStrategy fieldMissingStrategy)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:480
↓ 4 callersMethodsetFormatErrorStrategy
(DirtyDataStrategy formatErrorStrategy)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:304
↓ 4 callersMethodsetFormatErrorStrategy
(DirtyDataStrategy formatErrorStrategy)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:475
↓ 4 callersMethodsetRunning
(boolean running)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RunningChecker.java:28
↓ 4 callersMethodstart
start inner consumer
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:33
↓ 3 callersMethodasSummaryString
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:201
↓ 3 callersMethodcall
(Callable<T> callable, String errorMsg)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RetryUtil.java:46
↓ 3 callersMethodcollect
(RowData physicalRow)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:425
↓ 3 callersMethodcommittedOffsets
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 callersMethodcreateDynamicTableSink
(Map<String, String> options)
src/test/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSinkFactoryTest.java:97
↓ 3 callersMethodflush
(boolean endOfInput)
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/RocketMQWriter.java:108
↓ 3 callersMethodflushSync
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:221
↓ 3 callersMethodgetAclRpcHook
()
src/test/java/org/apache/flink/connector/rocketmq/example/ConnectorConfig.java:50
↓ 3 callersMethodgetCheckpoint
()
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:45
↓ 3 callersMethodgetData
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/BytesMessage.java:31
↓ 3 callersMethodgetHeaderValue
(BytesMessage message, int index)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:247
↓ 3 callersMethodgetIncreaseSet
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:679
↓ 3 callersMethodhashCode
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:154
↓ 3 callersMethodmaxOffsets
List max offsets for the specified MessageQueues.
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelector.java:75
↓ 3 callersMethodopen
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 callersMethodseek
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 callersMethodsendSplitChangesToRemote
(Set<Integer> pendingReaders)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:541
↓ 3 callersMethodsetProperties
(Map<String, String> props)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/BytesMessage.java:43
↓ 3 callersMethodsetTableSchema
(TableSchema tableSchema)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:470
↓ 3 callersMethodsetTopic
(String topic)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:61
↓ 3 callersMethodtoInt
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 callersMethodaddFinishedSplit
(String splitId)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:297
↓ 2 callersMethodallocate
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 callersMethodallocate
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 callersMethodassign
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 callersMethodbuild
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:573
↓ 2 callersMethodbuildAclRPCHook
Build credentials for client. @param props @return
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQConfig.java:183
↓ 2 callersMethodbuildCommonConfigs
Build Common Configs. @param props Properties @param client ClientConfig
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQConfig.java:166
↓ 2 callersMethodbuilder
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