MCPcopy Create free account

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

Functions733 in github.com/apache/rocketmq-flink

MethodgetOffsetsToCommit
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:241
MethodgetParallelInstanceId
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
MethodgetParallelInstanceId
()
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContextImpl.java:43
MethodgetProducedType
()
src/test/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceTest.java:57
MethodgetProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:654
MethodgetProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleStringDeserializationSchema.java:35
MethodgetProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleTupleDeserializationSchema.java:36
MethodgetProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:120
MethodgetProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/ForwardMessageExtDeserialization.java:33
MethodgetProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleKeyValueDeserializationSchema.java:63
MethodgetProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:189
MethodgetProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:395
MethodgetProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQSchemaWrapper.java:30
MethodgetProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchemaWrapper.java:53
MethodgetProducedType
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQRowDeserializationSchema.java:81
MethodgetProducerGroup
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
MethodgetProducerGroup
()
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:154
MethodgetProperties
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:138
MethodgetProperties
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/BytesMessage.java:39
MethodgetProperties
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
MethodgetQueueId
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:73
MethodgetQueueId
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:93
MethodgetQueueId
()
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:49
MethodgetQueueOffset
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:81
MethodgetQueueOffset
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:98
MethodgetReSendInitAssign
()
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceDetectEvent.java:27
MethodgetScanRuntimeProvider
(ScanContext scanContext)
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:119
MethodgetSinkRuntimeProvider
( DynamicTableSink.Context context)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:157
MethodgetSplitSerializer
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:179
MethodgetStoreSize
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:113
MethodgetStrategyName
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/ConsistentHashAllocateStrategy.java:29
MethodgetStrategyName
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AverageAllocateStrategy.java:12
MethodgetStrategyName
Allocate strategy name @return Current strategy name
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AllocateStrategy.java:36
MethodgetStrategyName
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/BroadcastAllocateStrategy.java:29
MethodgetTag
(T tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/TopicSelector.java:25
MethodgetTag
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
MethodgetTag
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:103
MethodgetTopic
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/selector/SimpleTopicSelectorTest.java:29
MethodgetTopic
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/selector/DefaultTopicSelectorTest.java:26
MethodgetTopic
(T tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/TopicSelector.java:23
MethodgetTopic
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:57
MethodgetTopic
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:83
MethodgetTopic
()
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:57
MethodgetUserCodeClassLoader
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:118
MethodgetValue
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:100
MethodgetVersion
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittableSerializer.java:33
MethodhandleSourceEvent
(int taskId, SourceEvent sourceEvent)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:251
MethodhandleSourceEvents
(SourceEvent sourceEvent)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:166
MethodhandleSourceQueueChange
(Set<MessageQueue> latestSet, Throwable t)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:300
MethodhandleSplitChanges
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
MethodhandleSplitRequest
(int subtaskId, @Nullable String requesterHostname)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:196
MethodinitializeState
(FunctionInitializationContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:237
MethodinitializeState
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
MethodinitializedState
(RocketMQSourceSplit partitionSplit)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:148
Methodinvoke
(Message input, Context context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:109
Methodinvoke
(RowData rowData, Context context)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataSink.java:49
MethodisEnableSchemaEvolution
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
MethodisEnableSchemaEvolution
()
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContextImpl.java:53
MethodisRunning
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceTest.java:109
Methodlatest
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
MethodlistReadableMetadata
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:150
MethodlistWritableMetadata
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:163
Methodmain
(String[] args)
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleConsumer.java:37
Methodmain
(String[] args)
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleProducer.java:40
Methodmain
(String[] args)
src/test/java/org/apache/flink/connector/rocketmq/example/LegacyConnectorExample.java:81
Methodmain
(String[] args)
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleAdmin.java:100
Methodmain
(String[] args)
src/test/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkTest.java:50
Methodmain
(String[] args)
src/test/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceTest.java:31
MethodmarkActive
()
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:86
MethodmarkIdle
()
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:82
MethodmaxOffsets
(Collection<MessageQueue> messageQueues)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:417
MethodminOffsets
(Collection<MessageQueue> messageQueues)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:388
MethodnextRecordFromSplit
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
MethodnextSplit
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
MethodnotifyAssignResult
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:641
MethodnotifyCheckpointComplete
(long checkpointId)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:659
MethodnotifyCheckpointComplete
(long checkpointId)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:130
Methodoffsets
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
MethodoffsetsForTimes
( Map<MessageQueue, Long> messageQueueWithTimeMap)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:446
MethodonException
(Throwable throwable)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:137
MethodonSplitFinished
(Map<String, RocketMQSourceSplitState> finishedSplitIds)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:93
MethodonSuccess
(SendResult sendResult)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:128
Methodopen
(Configuration parameters)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:77
Methodopen
(Configuration parameters)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:147
Methodopen
(Configuration configuration)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataSink.java:38
Methodopen
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataConverter.java:102
Methodopen
(InitializationContext context)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:123
Methodopen
(DeserializationSchema.InitializationContext context)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchemaWrapper.java:45
Methodopen
(InitializationContext context)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQRowDeserializationSchema.java:63
Methodopen
()
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:108
MethodoptionalOptions
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSinkFactory.java:73
MethodoptionalOptions
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQDynamicTableSourceFactory.java:65
Methodoverride
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
Methodpause
(Collection<MessageQueue> messageQueues)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:223
Methodpoll
(Duration timeout)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:189
MethodprepareCommit
()
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/RocketMQWriter.java:113
MethodprocessElement
(Tuple2<String, String> tuple, Context ctx, Collector<Message> out)
src/main/java/org/apache/flink/connector/rocketmq/legacy/function/SinkMapFunction.java:40
MethodprocessElement
( Tuple2<String, String> value, Context ctx, Collector<Tuple2<String, String>> out)
src/main/java/org/apache/flink/connector/rocketmq/legacy/function/SourceMapFunction.java:27
MethodprocessTime
Returns the current process time in flink.
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContext.java:51
MethodprocessTime
()
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContextImpl.java:58
← previousnext →501–600 of 733, ranked by callers