MCPcopy Create free account

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

Functions733 in github.com/apache/rocketmq-flink

Methodbuild
()
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigValidator.java:99
MethodcanNotSetSameOptionTwiceWithDifferentValue
()
src/test/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigBuilderTest.java:17
MethodcheckAndGetNextWatermark
(MessageExt lastElement, long extractedTimestamp)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/PunctuatedAssigner.java:43
MethodcheckLocalTransaction
(MessageExt msg)
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:115
Methodclone
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/RocketMQCommitter.java:102
Methodclone
(RocketMQSourceSplit split)
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:132
Methodclose
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSinkTest.java:67
Methodclose
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceTest.java:101
Methodclose
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:208
Methodclose
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:529
Methodclose
()
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:302
Methodclose
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/RocketMQCommitter.java:95
Methodclose
()
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/RocketMQWriter.java:133
Methodclose
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataSink.java:57
Methodclose
()
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:354
Methodclose
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitter.java:66
Methodclose
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:229
Methodclose
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:446
Methodclose
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:239
Methodclose
()
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:126
Methodcommit
(SendCommittable sendCommittable)
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:288
Methodcommit
(Collection<CommitRequest<SendCommittable>> requests)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/RocketMQCommitter.java:51
MethodcommitOffset
(MessageQueue messageQueue, long offset)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:336
MethodcommittedOffsets
(Collection<MessageQueue> messageQueues)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:359
MethodcommittedOffsets
The group id should be the set for {@link RocketMQSource } before invoking this method. Otherwise, an {@code IllegalStateException} will be thrown.
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelector.java:69
MethodconflictOptions
(ConfigOption<?>... options)
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigValidator.java:88
MethodconsistentHashAllocateStrategyTest
()
src/test/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/ConsistentHashAllocateStrategyTest.java:21
Methodcopy
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:176
Methodcopy
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:163
MethodcreateCommitter
()
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSink.java:64
MethodcreateDynamicTableSink
(Context context)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSinkFactory.java:91
MethodcreateDynamicTableSource
(Context context)
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQDynamicTableSourceFactory.java:86
MethodcreateEnumerator
( SplitEnumeratorContext<RocketMQSourceSplit> enumContext)
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:153
MethodcreateOutputForSplit
(String splitId)
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:89
MethodcreateReader
(SourceReaderContext readerContext)
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:104
MethodcreateRocketMQDeserializationSchema
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:193
MethodcreateWriter
(InitContext context)
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSink.java:58
Methoddeserialize
( MessageView messageView, Collector<String> out)
src/test/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceTest.java:51
Methoddeserialize
(int version, byte[] serialized)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittableSerializer.java:54
Methoddeserialize
Deserializes the byte message. <p>Can output multiple records through the {@link Collector}. Note that number and size of the produced records should
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchema.java:59
Methoddeserialize
(List<BytesMessage> messages, Collector<RowData> collector)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:133
Methoddeserialize
(MessageView messageView, Collector<T> out)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQSchemaWrapper.java:27
Methoddeserialize
(MessageView messageView, Collector<T> out)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchemaWrapper.java:50
Methoddeserialize
(MessageView messageView, Collector<RowData> collector)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQRowDeserializationSchema.java:69
Methoddeserialize
(byte[] value, ValueType type)
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteSerializer.java:40
MethoddeserializeKeyAndValue
(byte[] key, byte[] value)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleStringDeserializationSchema.java:29
MethoddeserializeKeyAndValue
(byte[] key, byte[] value)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleTupleDeserializationSchema.java:29
MethoddeserializeMessageExt
Deserialize messageExt to type T you want to output. @param messageExt the messageExt @return the t
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/MessageExtDeserializationScheme.java:38
MethoddeserializeMessageExt
(MessageExt messageExt)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/ForwardMessageExtDeserialization.java:28
MethodexecuteLocalTransaction
(Message msg, Object arg)
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:108
MethodextractMessages
(List<MessageExt> messages)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQRowDeserializationSchema.java:86
MethodextractTimestamp
(MessageQueue mq, long timestamp)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/WaterMarkPerQueue.java:41
MethodextractTimestamp
(MessageExt element, long previousElementTimestamp)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/BoundedOutOfOrdernessGenerator.java:37
MethodextractTimestamp
(MessageExt element, long previousElementTimestamp)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/PunctuatedAssigner.java:38
MethodextractTimestamp
(MessageExt element, long previousElementTimestamp)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/BoundedOutOfOrdernessGeneratorPerQueue.java:41
MethodextractTimestamp
(MessageExt element, long previousElementTimestamp)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/TimeLagWatermarkGenerator.java:38
MethodfactoryIdentifier
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSinkFactory.java:57
MethodfactoryIdentifier
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQDynamicTableSourceFactory.java:52
Methodfetch
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:116
MethodfetchMessageQueues
(String topic)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:152
MethodfinishedSplits
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:342
MethodflinkSchema
Create a RocketMQDeserializationSchema by using the flink's {@link DeserializationSchema}. It would consume the rocketmq message as byte array and dec
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchema.java:66
MethodgetAutoOffsetResetStrategy
Returns the strategy for automatically resetting the offset when there is no initial offset in RocketMQ or if the current offset does not exist in Roc
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelector.java:57
MethodgetAutoOffsetResetStrategy
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorByStrategy.java:56
MethodgetAutoOffsetResetStrategy
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorByTimestamp.java:63
MethodgetAutoOffsetResetStrategy
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorNoStopping.java:39
MethodgetAutoOffsetResetStrategy
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorBySpecified.java:77
MethodgetBody
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:118
MethodgetBoundedness
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:99
MethodgetBrokerName
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:65
MethodgetBrokerName
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:88
MethodgetBytesMessages
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQRowDeserializationSchema.java:108
MethodgetChangelogMode
(ChangelogMode requestedMode)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:146
MethodgetChangelogMode
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:114
MethodgetCommittableSerializer
()
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSink.java:69
MethodgetConsumerGroup
()
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:147
MethodgetCurrentWatermark
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/WaterMarkPerQueue.java:46
MethodgetCurrentWatermark
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/BoundedOutOfOrdernessGenerator.java:44
MethodgetCurrentWatermark
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/BoundedOutOfOrdernessGeneratorPerQueue.java:50
MethodgetCurrentWatermark
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/TimeLagWatermarkGenerator.java:43
MethodgetDecreaseSet
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:713
MethodgetDeliveryAttempt
Get the number of times that the message has been attempted to be delivered. @return the delivery attempt count
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:94
MethodgetDeliveryAttempt
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:123
MethodgetDescription
()
src/main/java/org/apache/flink/connector/rocketmq/source/config/OffsetVerification.java:44
MethodgetEnumeratorCheckpointSerializer
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:184
MethodgetEventTime
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:128
MethodgetIncreaseSet
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:709
MethodgetIngestionTime
Get the ingestion time of the message, which is the time that the message was received by the broker. @return the ingestion time
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:109
MethodgetIngestionTime
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:133
MethodgetKeys
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:108
MethodgetMailboxExecutor
()
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContextImpl.java:63
MethodgetMessageId
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:78
MethodgetMessageQueue
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:116
MethodgetMessageQueueOffsets
( Collection<MessageQueue> messageQueues, MessageQueueOffsetsRetriever offsetsRetriever)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorByStrategy.java:40
MethodgetMessageQueueOffsets
( Collection<MessageQueue> messageQueues, MessageQueueOffsetsRetriever offsetsRetriever)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorByTimestamp.java:36
MethodgetMessageQueueOffsets
( Collection<MessageQueue> messageQueues, MessageQueueOffsetsRetriever offsetsRetriever)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorNoStopping.java:32
MethodgetMessageQueueOffsets
( Collection<MessageQueue> messageQueues, MessageQueueOffsetsRetriever offsetsRetriever)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorBySpecified.java:43
MethodgetNumAliveFetchers
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:246
MethodgetNumberOfParallelInstances
@return number of parallel RocketMQSink tasks.
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContext.java:41
MethodgetNumberOfParallelInstances
()
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContextImpl.java:48
← previousnext →401–500 of 733, ranked by callers