MCPcopy Create free account

hub / github.com/apache/rocketmq-flink / types & classes

Types & classes161 in github.com/apache/rocketmq-flink

InterfaceAllocateStrategy
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
ClassAllocateStrategyFactory
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AllocateStrategyFactory.java:26
ClassAverageAllocateStrategy
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AverageAllocateStrategy.java:11
ClassAverageAllocateStrategyTest
src/test/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AverageAllocateStrategyTest.java:15
ClassBoundedOutOfOrdernessGenerator
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/BoundedOutOfOrdernessGenerator.java:25
ClassBoundedOutOfOrdernessGeneratorPerQueue
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/BoundedOutOfOrdernessGeneratorPerQueue.java:28
ClassBroadcastAllocateStrategy
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/BroadcastAllocateStrategy.java:27
ClassBroadcastAllocateStrategyTest
src/test/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/BroadcastAllocateStrategyTest.java:15
ClassBuilder
Builder of {@link RowKeyValueDeserializationSchema}.
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:286
ClassBuilder
Builder of {@link RowDeserializationSchema}.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:453
ClassByteSerializer
BytesSerializer is responsible to deserialize field from byte array.
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteSerializer.java:36
ClassByteUtils
Utility class to for operations related to bytes.
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteUtils.java:25
ClassBytesMessage
Message contains byte array.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/BytesMessage.java:25
ClassCollectorOption
Options for {@link RowKeyValueDeserializationSchema}.
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:395
ClassCollectorOption
Options for {@link RowDeserializationSchema}.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:591
ClassConnectorConfig
src/test/java/org/apache/flink/connector/rocketmq/example/ConnectorConfig.java:24
ClassConsistentHashAllocateStrategy
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/ConsistentHashAllocateStrategy.java:27
ClassConsistentHashAllocateStrategyTest
src/test/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/ConsistentHashAllocateStrategyTest.java:14
ClassDefaultTopicSelector
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/DefaultTopicSelector.java:19
ClassDefaultTopicSelectorTest
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/selector/DefaultTopicSelectorTest.java:25
EnumDirtyDataStrategy
Dirty data process strategy.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/DirtyDataStrategy.java:21
ClassForwardMessageExtDeserialization
A Forward messageExt deserialization.
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/ForwardMessageExtDeserialization.java:25
ClassHashMessageQueueSelector
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/HashMessageQueueSelector.java:26
ClassHashMessageQueueSelectorTest
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/selector/HashMessageQueueSelectorTest.java:30
InterfaceInnerConsumer
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:30
ClassInnerConsumerImpl
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:59
InterfaceInnerProducer
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
ClassInnerProducerImpl
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:54
InterfaceKeyValueDeserializationSchema
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/KeyValueDeserializationSchema.java:23
InterfaceKeyValueSerializationSchema
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/KeyValueSerializationSchema.java:21
ClassLatencyGauge
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:78
ClassLegacyConnectorExample
src/test/java/org/apache/flink/connector/rocketmq/example/LegacyConnectorExample.java:42
InterfaceMessageExtDeserializationScheme
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
InterfaceMessageQueueOffsetsRetriever
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
InterfaceMessageQueueSelector
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/MessageQueueSelector.java:23
InterfaceMessageView
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
ClassMessageViewExt
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:27
ClassMetadataCollector
Metadata of RowData collector.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:410
InterfaceMetadataConverter
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataConverter.java:213
InterfaceMetadataConverter
Source metadata converter interface.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:403
ClassMetricUtil
src/main/java/org/apache/flink/connector/rocketmq/MetricUtil.java:20
ClassMetricUtils
RocketMQ connector metrics.
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:29
EnumOffsetResetStrategy
Config for #{@link StartupMode#GROUP_OFFSETS}.
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/config/OffsetResetStrategy.java:21
EnumOffsetVerification
The enum class for defining the offset verify behavior.
src/main/java/org/apache/flink/connector/rocketmq/source/config/OffsetVerification.java:28
InterfaceOffsetsSelector
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
ClassOffsetsSelectorBySpecified
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorBySpecified.java:32
ClassOffsetsSelectorByStrategy
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorByStrategy.java:29
ClassOffsetsSelectorByTimestamp
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorByTimestamp.java:28
ClassOffsetsSelectorNoStopping
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorNoStopping.java:29
InterfaceOffsetsValidator
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsValidator.java:7
ClassPunctuatedAssigner
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
InterfaceQueryableSchema
An interface for the deserialization of records.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/QueryableSchema.java:30
ClassRandomMessageQueueSelector
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/RandomMessageQueueSelector.java:27
ClassRandomMessageQueueSelectorTest
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/selector/RandomMessageQueueSelectorTest.java:30
EnumReadableMetadata
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:237
ClassRemotingOffsetsRetrieverImpl
The implementation for offsets retriever with a consumer and an admin client.
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:344
ClassRetryUtil
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RetryUtil.java:26
ClassRetryUtilTest
Tests for {@link RetryUtil}.
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/util/RetryUtilTest.java:31
ClassRocketMQCatalog
A catalog implementation for RocketMQ.
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:80
ClassRocketMQCatalogFactory
The {@CatalogFactory} implementation of RocketMQ.
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactory.java:36
ClassRocketMQCatalogFactoryOptions
{@link ConfigOption}s for {@link RocketMQCatalog}.
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactoryOptions.java:29
ClassRocketMQCatalogFactoryTest
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactoryTest.java:33
ClassRocketMQCatalogTest
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:66
ClassRocketMQCommitter
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
ClassRocketMQConfig
RocketMQConfig for Consumer/Producer.
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQConfig.java:36
ClassRocketMQConfigBuilder
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
ClassRocketMQConfigBuilderTest
src/test/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigBuilderTest.java:15
ClassRocketMQConfigValidator
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
ClassRocketMQConfigValidatorBuilder
Builder pattern for building {@link RocketMQConfigValidator}.
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigValidator.java:83
ClassRocketMQConfiguration
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
InterfaceRocketMQDeserializationSchema
An interface for the deserialization of RocketMQ records.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchema.java:32
ClassRocketMQDeserializationSchemaWrapper
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
ClassRocketMQDynamicTableSink
Defines the dynamic table sink of RocketMQ.
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:44
ClassRocketMQDynamicTableSinkFactory
Defines the {@link DynamicTableSinkFactory} implementation to create {@link RocketMQDynamicTableSink}.
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSinkFactory.java:55
ClassRocketMQDynamicTableSinkFactoryTest
Tests for {@link RocketMQDynamicTableSinkFactory}.
src/test/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSinkFactoryTest.java:46
ClassRocketMQDynamicTableSourceFactory
Defines the {@link DynamicTableSourceFactory} implementation to create {@link RocketMQScanTableSource}.
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQDynamicTableSourceFactory.java:48
ClassRocketMQDynamicTableSourceFactoryTest
Tests for {@link RocketMQDynamicTableSourceFactory}.
src/test/java/org/apache/flink/connector/rocketmq/source/table/RocketMQDynamicTableSourceFactoryTest.java:46
ClassRocketMQOptions
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
ClassRocketMQPartitionSplitSerializer
The {@link SimpleVersionedSerializer serializer} for {@link RocketMQSourceSplit}.
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQPartitionSplitSerializer.java:30
ClassRocketMQPartitionSplitSerializerTest
Test for {@link RocketMQPartitionSplitSerializer}.
src/test/java/org/apache/flink/connector/rocketmq/source/split/RocketMQPartitionSplitSerializerTest.java:28
ClassRocketMQRecordEmitter
The {@link RecordEmitter} implementation for {@link RocketMQSourceReader}.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitter.java:30
ClassRocketMQRecordEmitterTest
Test for {@link RocketMQRecordEmitter}.
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:37
ClassRocketMQRecordsWithSplitIds
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:275
ClassRocketMQRowDataConverter
RocketMQRowDataConverter converts the row data of table to RocketMQ message pattern.
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataConverter.java:44
ClassRocketMQRowDataSink
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
ClassRocketMQRowDeserializationSchema
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
ClassRocketMQRowDeserializationSchemaTest
Test for {@link RocketMQRowDeserializationSchema}.
src/test/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQRowDeserializationSchemaTest.java:24
ClassRocketMQScanTableSource
Defines the scan table source of RocketMQ.
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:50
ClassRocketMQSchemaWrapper
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQSchemaWrapper.java:26
InterfaceRocketMQSerializationSchema
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
ClassRocketMQSerializerWrapper
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
ClassRocketMQSink
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
ClassRocketMQSink
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSink.java:32
ClassRocketMQSinkBuilder
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
InterfaceRocketMQSinkContext
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
ClassRocketMQSinkContextImpl
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContextImpl.java:26
ClassRocketMQSinkOptions
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkOptions.java:27
ClassRocketMQSinkTest
src/test/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSinkTest.java:39
ClassRocketMQSinkTest
src/test/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkTest.java:46
ClassRocketMQSource
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