Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/apache/rocketmq-flink
/ types & classes
Types & classes
161 in github.com/apache/rocketmq-flink
⨍
Functions
733
◇
Types & classes
161
Interface
AllocateStrategy
This interface defines a strategy for allocating RocketMQ source splits to Flink tasks.
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AllocateStrategy.java:28
Class
AllocateStrategyFactory
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AllocateStrategyFactory.java:26
Class
AverageAllocateStrategy
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AverageAllocateStrategy.java:11
Class
AverageAllocateStrategyTest
src/test/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AverageAllocateStrategyTest.java:15
Class
BoundedOutOfOrdernessGenerator
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/BoundedOutOfOrdernessGenerator.java:25
Class
BoundedOutOfOrdernessGeneratorPerQueue
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/BoundedOutOfOrdernessGeneratorPerQueue.java:28
Class
BroadcastAllocateStrategy
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/BroadcastAllocateStrategy.java:27
Class
BroadcastAllocateStrategyTest
src/test/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/BroadcastAllocateStrategyTest.java:15
Class
Builder
Builder of {@link RowKeyValueDeserializationSchema}.
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:286
Class
Builder
Builder of {@link RowDeserializationSchema}.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:453
Class
ByteSerializer
BytesSerializer is responsible to deserialize field from byte array.
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteSerializer.java:36
Class
ByteUtils
Utility class to for operations related to bytes.
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteUtils.java:25
Class
BytesMessage
Message contains byte array.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/BytesMessage.java:25
Class
CollectorOption
Options for {@link RowKeyValueDeserializationSchema}.
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:395
Class
CollectorOption
Options for {@link RowDeserializationSchema}.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:591
Class
ConnectorConfig
src/test/java/org/apache/flink/connector/rocketmq/example/ConnectorConfig.java:24
Class
ConsistentHashAllocateStrategy
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/ConsistentHashAllocateStrategy.java:27
Class
ConsistentHashAllocateStrategyTest
src/test/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/ConsistentHashAllocateStrategyTest.java:14
Class
DefaultTopicSelector
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/DefaultTopicSelector.java:19
Class
DefaultTopicSelectorTest
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/selector/DefaultTopicSelectorTest.java:25
Enum
DirtyDataStrategy
Dirty data process strategy.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/DirtyDataStrategy.java:21
Class
ForwardMessageExtDeserialization
A Forward messageExt deserialization.
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/ForwardMessageExtDeserialization.java:25
Class
HashMessageQueueSelector
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/HashMessageQueueSelector.java:26
Class
HashMessageQueueSelectorTest
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/selector/HashMessageQueueSelectorTest.java:30
Interface
InnerConsumer
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:30
Class
InnerConsumerImpl
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:59
Interface
InnerProducer
InnerProducer is an interface that represents a message producer used for sending messages to a messaging system. @see AutoCloseable
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducer.java:33
Class
InnerProducerImpl
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:54
Interface
KeyValueDeserializationSchema
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/KeyValueDeserializationSchema.java:23
Interface
KeyValueSerializationSchema
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/KeyValueSerializationSchema.java:21
Class
LatencyGauge
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:78
Class
LegacyConnectorExample
src/test/java/org/apache/flink/connector/rocketmq/example/LegacyConnectorExample.java:42
Interface
MessageExtDeserializationScheme
The interface Message ext deserialization scheme. @param <T> the type parameter
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/MessageExtDeserializationScheme.java:31
Interface
MessageQueueOffsetsRetriever
An interface that provides necessary information to the {@link OffsetsSelector} to get the initial offsets of the RocketMQ message queues.
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelector.java:63
Interface
MessageQueueSelector
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/MessageQueueSelector.java:23
Interface
MessageView
This interface defines the methods for obtaining information about a message in RocketMQ.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:24
Class
MessageViewExt
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:27
Class
MetadataCollector
Metadata of RowData collector.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:410
Interface
MetadataConverter
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataConverter.java:213
Interface
MetadataConverter
Source metadata converter interface.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:403
Class
MetricUtil
src/main/java/org/apache/flink/connector/rocketmq/MetricUtil.java:20
Class
MetricUtils
RocketMQ connector metrics.
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:29
Enum
OffsetResetStrategy
Config for #{@link StartupMode#GROUP_OFFSETS}.
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/config/OffsetResetStrategy.java:21
Enum
OffsetVerification
The enum class for defining the offset verify behavior.
src/main/java/org/apache/flink/connector/rocketmq/source/config/OffsetVerification.java:28
Interface
OffsetsSelector
An interface for users to specify the starting / stopping offset of a {@link RocketMQSourceSplit}.
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelector.java:37
Class
OffsetsSelectorBySpecified
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorBySpecified.java:32
Class
OffsetsSelectorByStrategy
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorByStrategy.java:29
Class
OffsetsSelectorByTimestamp
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorByTimestamp.java:28
Class
OffsetsSelectorNoStopping
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorNoStopping.java:29
Interface
OffsetsValidator
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsValidator.java:7
Class
PunctuatedAssigner
With Punctuated Watermarks To generate watermarks whenever a certain event indicates that a new watermark might be generated, use AssignerWithPunctuat
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/PunctuatedAssigner.java:37
Interface
QueryableSchema
An interface for the deserialization of records.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/QueryableSchema.java:30
Class
RandomMessageQueueSelector
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/RandomMessageQueueSelector.java:27
Class
RandomMessageQueueSelectorTest
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/selector/RandomMessageQueueSelectorTest.java:30
Enum
ReadableMetadata
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:237
Class
RemotingOffsetsRetrieverImpl
The implementation for offsets retriever with a consumer and an admin client.
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:344
Class
RetryUtil
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RetryUtil.java:26
Class
RetryUtilTest
Tests for {@link RetryUtil}.
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/util/RetryUtilTest.java:31
Class
RocketMQCatalog
A catalog implementation for RocketMQ.
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:80
Class
RocketMQCatalogFactory
The {@CatalogFactory} implementation of RocketMQ.
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactory.java:36
Class
RocketMQCatalogFactoryOptions
{@link ConfigOption}s for {@link RocketMQCatalog}.
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactoryOptions.java:29
Class
RocketMQCatalogFactoryTest
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactoryTest.java:33
Class
RocketMQCatalogTest
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:66
Class
RocketMQCommitter
Committer implementation for {@link RocketMQSink} <p>The committer is responsible to finalize the RocketMQ transactions by committing them.
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/RocketMQCommitter.java:41
Class
RocketMQConfig
RocketMQConfig for Consumer/Producer.
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQConfig.java:36
Class
RocketMQConfigBuilder
A builder for building the unmodifiable {@link Configuration} instance. Providing the common validate logic for RocketMQ source & sink.
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigBuilder.java:40
Class
RocketMQConfigBuilderTest
src/test/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigBuilderTest.java:15
Class
RocketMQConfigValidator
A config validator for building {@link RocketMQConfiguration} in {@link RocketMQConfigBuilder}. It's used for source & sink builder. <p>We would vali
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigValidator.java:45
Class
RocketMQConfigValidatorBuilder
Builder pattern for building {@link RocketMQConfigValidator}.
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigValidator.java:83
Class
RocketMQConfiguration
An unmodifiable {@link Configuration} for RocketMQ. We provide extra methods for building the different RocketMQ client instance.
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfiguration.java:38
Interface
RocketMQDeserializationSchema
An interface for the deserialization of RocketMQ records.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchema.java:32
Class
RocketMQDeserializationSchemaWrapper
A {@link RocketMQDeserializationSchema} implementation which based on the given flink's {@link DeserializationSchema}. We would consume the message as
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchemaWrapper.java:35
Class
RocketMQDynamicTableSink
Defines the dynamic table sink of RocketMQ.
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:44
Class
RocketMQDynamicTableSinkFactory
Defines the {@link DynamicTableSinkFactory} implementation to create {@link RocketMQDynamicTableSink}.
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSinkFactory.java:55
Class
RocketMQDynamicTableSinkFactoryTest
Tests for {@link RocketMQDynamicTableSinkFactory}.
src/test/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSinkFactoryTest.java:46
Class
RocketMQDynamicTableSourceFactory
Defines the {@link DynamicTableSourceFactory} implementation to create {@link RocketMQScanTableSource}.
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQDynamicTableSourceFactory.java:48
Class
RocketMQDynamicTableSourceFactoryTest
Tests for {@link RocketMQDynamicTableSourceFactory}.
src/test/java/org/apache/flink/connector/rocketmq/source/table/RocketMQDynamicTableSourceFactoryTest.java:46
Class
RocketMQOptions
Configuration for RocketMQ Client, these config options would be used for both source, sink and table.
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQOptions.java:35
Class
RocketMQPartitionSplitSerializer
The {@link SimpleVersionedSerializer serializer} for {@link RocketMQSourceSplit}.
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQPartitionSplitSerializer.java:30
Class
RocketMQPartitionSplitSerializerTest
Test for {@link RocketMQPartitionSplitSerializer}.
src/test/java/org/apache/flink/connector/rocketmq/source/split/RocketMQPartitionSplitSerializerTest.java:28
Class
RocketMQRecordEmitter
The {@link RecordEmitter} implementation for {@link RocketMQSourceReader}.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitter.java:30
Class
RocketMQRecordEmitterTest
Test for {@link RocketMQRecordEmitter}.
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:37
Class
RocketMQRecordsWithSplitIds
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:275
Class
RocketMQRowDataConverter
RocketMQRowDataConverter converts the row data of table to RocketMQ message pattern.
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataConverter.java:44
Class
RocketMQRowDataSink
RocketMQRowDataSink helps for writing the converted row data of table to RocketMQ messages.
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataSink.java:26
Class
RocketMQRowDeserializationSchema
A row data wrapper class that wraps a {@link RocketMQDeserializationSchema} to deserialize {@link MessageExt}.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQRowDeserializationSchema.java:41
Class
RocketMQRowDeserializationSchemaTest
Test for {@link RocketMQRowDeserializationSchema}.
src/test/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQRowDeserializationSchemaTest.java:24
Class
RocketMQScanTableSource
Defines the scan table source of RocketMQ.
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:50
Class
RocketMQSchemaWrapper
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQSchemaWrapper.java:26
Interface
RocketMQSerializationSchema
The serialization schema for how to serialize records into RocketMQ. A serialization schema which defines how to convert a value of type {@code T} to
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/serializer/RocketMQSerializationSchema.java:18
Class
RocketMQSerializerWrapper
Wrap the RocketMQ Schema into RocketMQSerializationSchema. We support schema evolution out of box by this implementation.
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/serializer/RocketMQSerializerWrapper.java:26
Class
RocketMQSink
The RocketMQSink provides at-least-once reliability guarantees when checkpoints are enabled and batchFlushOnCheckpoint(true) is set. Otherwise, the si
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:51
Class
RocketMQSink
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSink.java:32
Class
RocketMQSinkBuilder
Builder to construct {@link RocketMQSink}. @see RocketMQSink for a more detailed explanation of the different guarantees.
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:46
Interface
RocketMQSinkContext
This context provides information on the rocketmq record target location. An implementation that would contain all the required context.
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContext.java:28
Class
RocketMQSinkContextImpl
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContextImpl.java:26
Class
RocketMQSinkOptions
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkOptions.java:27
Class
RocketMQSinkTest
src/test/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSinkTest.java:39
Class
RocketMQSinkTest
src/test/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkTest.java:46
Class
RocketMQSource
The Source implementation of RocketMQ.
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:57
next →
1–100 of 161, ranked by callers