MCPcopy Create free account

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

Functions733 in github.com/apache/rocketmq-flink

↓ 1 callersMethodsetInitFinish
(boolean initFinish)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:258
↓ 1 callersMethodsetLineDelimiter
(String lineDelimiter)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:500
↓ 1 callersMethodsetMetadataConverters
(MetadataConverter[] metadataConverters)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:515
↓ 1 callersMethodsetOffsetMsgId
(String offsetMsgId)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:101
↓ 1 callersMethodsetProperties
Set arbitrary properties for the RocketMQSink and RocketMQ Consumer. The valid keys can be found in {@link RocketMQSinkOptions} and {@link RocketMQOpt
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:136
↓ 1 callersMethodsetQueueId
(Integer queueId)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:77
↓ 1 callersMethodsetQueueId
(int queueId)
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:37
↓ 1 callersMethodsetQueueOffset
(Long queueOffset)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:85
↓ 1 callersMethodsetSerializer
Sets the {@link RocketMQSerializationSchema} that transforms incoming records to {@link org.apache.rocketmq.common.message.MessageExt}s. @param seria
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:148
↓ 1 callersMethodsetSourceOutput
(SourceOutput<T> sourceOutput)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitter.java:69
↓ 1 callersMethodsetStartFromEarliest
consume from the min offset at every restart with no state
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:442
↓ 1 callersMethodsetStartFromSpecificOffsets
consume from the specific offset. Group offsets is enable while the broker didn't specify offset.
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:481
↓ 1 callersMethodsetStartFromTimeStamp
consume from the closest offset
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:454
↓ 1 callersMethodsetTimestamp
(long timestamp)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitter.java:73
↓ 1 callersMethodsetTopic
(String topic)
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:41
↓ 1 callersMethodsetTopics
Set a list of topics the RocketMQSource should consume from. All the topics in the list should have existed in the RocketMQ cluster. Otherwise, an exc
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:93
↓ 1 callersMethodsetUnbounded
(OffsetsSelector stoppingOffsetsSelector)
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:116
↓ 1 callersMethodsplitId
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:142
↓ 1 callersMethodtoBigDecimal
(byte[] bytes)
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteUtils.java:205
↓ 1 callersMethodtoBoolean
Convert a byte array to a boolean. @param b array @return True or false.
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteUtils.java:64
↓ 1 callersMethodtoDouble
Parse a byte array to double. @param bytes byte array @return Return double made from passed bytes.
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteUtils.java:140
↓ 1 callersMethodtoFloat
Presumes float encoded as IEEE 754 floating-point "single format". @param bytes byte array @return Float made from passed byte array.
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteUtils.java:119
↓ 1 callersMethodtoShort
Converts a byte array to a short value. @param bytes byte array @return the short value
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteUtils.java:161
↓ 1 callersMethodtoSplitId
(MessageQueue messageQueue)
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:124
↓ 1 callersMethodunregisterMessageQueue
(MessageQueue messageQueue)
src/main/java/org/apache/flink/connector/rocketmq/source/metrics/RocketMQSourceReaderMetrics.java:49
↓ 1 callersMethodupdateMessageQueueOffset
(MessageQueue mq, long offset)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:488
↓ 1 callersMethodwithAsync
(boolean async)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:183
↓ 1 callersMethodwithBatchFlushOnCheckpoint
(boolean batchFlushOnCheckpoint)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:188
↓ 1 callersMethodwithBatchSize
(int batchSize)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:193
MethodAllocateStrategyFactory
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AllocateStrategyFactory.java:32
MethodBoundedOutOfOrdernessGenerator
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/BoundedOutOfOrdernessGenerator.java:31
MethodBoundedOutOfOrdernessGeneratorPerQueue
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/BoundedOutOfOrdernessGeneratorPerQueue.java:34
MethodBuilder
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:297
MethodBuilder
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:468
MethodDefaultTopicSelector
(final String topicName, final String tagName)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/DefaultTopicSelector.java:23
MethodInnerConsumerImpl
(Configuration configuration)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:68
MethodInnerProducerImpl
(Configuration configuration)
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:65
MethodMessageViewExt
(MessageExt messageExt)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageViewExt.java:59
MethodMetadataCollector
(boolean hasMetadata, MetadataConverter[] metadataConverters)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:420
MethodOffsetVerification
(InlineElement desc)
src/main/java/org/apache/flink/connector/rocketmq/source/config/OffsetVerification.java:40
MethodOffsetsSelectorBySpecified
( Map<MessageQueue, Long> initialOffsets, OffsetResetStrategy offsetResetStrategy)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorBySpecified.java:37
MethodOffsetsSelectorByStrategy
( ConsumeFromWhere consumeFromWhere, OffsetResetStrategy offsetResetStrategy)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorByStrategy.java:34
MethodOffsetsSelectorByTimestamp
(long startingTimestamp)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelectorByTimestamp.java:32
MethodReadableMetadata
(String key, DataType dataType, MetadataConverter converter)
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:257
MethodRemotingOffsetsRetrieverImpl
(InnerConsumer innerConsumer)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:350
MethodRetryUtil
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RetryUtil.java:36
MethodRocketMQCatalog
( String catalogName, String database, String namesrvAddr, String schemaRegistryUrl)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:89
MethodRocketMQCatalogFactoryOptions
()
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactoryOptions.java:51
MethodRocketMQCommitter
(Configuration configuration)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/RocketMQCommitter.java:47
MethodRocketMQConfigValidator
( List<Set<ConfigOption<?>>> conflictOptions, Set<ConfigOption<?>> requiredOptions)
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigValidator.java:51
MethodRocketMQConfiguration
Creates a new RocketMQConfiguration, which holds a copy of the given configuration that can't be altered. @param config The configuration with the or
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfiguration.java:48
MethodRocketMQDeserializationSchemaWrapper
(DeserializationSchema<T> deserializationSchema)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchemaWrapper.java:41
MethodRocketMQDynamicTableSink
( DescriptorProperties properties, TableSchema schema, String topic,
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:71
MethodRocketMQRecordEmitter
(RocketMQDeserializationSchema<T> deserializationSchema)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitter.java:36
MethodRocketMQRecordsWithSplitIds
(RocketMQSourceReaderMetrics readerMetrics)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:288
MethodRocketMQRowDataConverter
( String topic, String tag, String dynamicColumn, String field
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataConverter.java:73
MethodRocketMQRowDataSink
(RocketMQSink sink, RocketMQRowDataConverter converter)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataSink.java:33
MethodRocketMQRowDeserializationSchema
( TableSchema tableSchema, Map<String, String> properties, boolean hasMeta
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQRowDeserializationSchema.java:49
MethodRocketMQScanTableSource
( long pollTime, DescriptorProperties properties, TableSchema schema,
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:76
MethodRocketMQSink
(Properties props)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:73
MethodRocketMQSink
( Configuration configuration, MessageQueueSelector messageQueueSelector,
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSink.java:39
MethodRocketMQSinkBuilder
()
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:58
MethodRocketMQSinkContextImpl
(InitContext initContext, Configuration configuration)
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContextImpl.java:35
MethodRocketMQSource
( OffsetsSelector startingOffsetsSelector, OffsetsSelector stoppingOffsetsSelector,
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:77
MethodRocketMQSourceBuilder
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:58
MethodRocketMQSourceEnumState
(Set<MessageQueue> currentSplitAssignment)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumState.java:33
MethodRocketMQSourceEnumerator
( OffsetsSelector startingOffsetsSelector, OffsetsSelector stoppingOffsetsSelector,
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:110
MethodRocketMQSourceFetcherManager
Creates a new SplitFetcherManager with a single I/O threads. @param elementsQueue The queue that is used to hand over data from the I/O thread (the
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceFetcherManager.java:55
MethodRocketMQSourceFunction
(KeyValueDeserializationSchema<OUT> schema, Properties props)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:142
MethodRocketMQSourceReader
( FutureCompletingBlockingQueue<RecordsWithSplitIds<MessageView>> elementsQueue, Rocke
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:75
MethodRocketMQSourceReaderMetrics
(SourceReaderMetricGroup sourceReaderMetricGroup)
src/main/java/org/apache/flink/connector/rocketmq/source/metrics/RocketMQSourceReaderMetrics.java:45
MethodRocketMQSourceSplit
( MessageQueue messageQueue, long startingOffset, long stoppingOffset)
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplit.java:43
MethodRocketMQSourceSplitState
(RocketMQSourceSplit partitionSplit)
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplitState.java:26
MethodRocketMQSplitReader
( Configuration configuration, SourceReaderContext sourceReaderContext, Ro
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:93
MethodRocketMQWriter
( Configuration configuration, MessageQueueSelector messageQueueSelector,
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/RocketMQWriter.java:66
MethodRowDeserializationSchema
( TableSchema tableSchema, DirtyDataStrategy formatErrorStrategy, DirtyDat
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:79
MethodRowKeyValueDeserializationSchema
( TableSchema tableSchema, DirtyDataStrategy formatErrorStrategy, DirtyDat
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:70
MethodSendCommittable
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:45
MethodSimpleAdmin
()
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleAdmin.java:41
MethodSimpleKeyValueDeserializationSchema
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleKeyValueDeserializationSchema.java:34
MethodSimpleKeyValueSerializationSchema
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleKeyValueSerializationSchema.java:29
MethodSimpleTopicSelector
SimpleTopicSelector Constructor. @param topicFieldName field name used for selecting topic @param defaultTopicName default field name used for select
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/SimpleTopicSelector.java:42
MethodSinkMapFunction
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/function/SinkMapFunction.java:33
MethodSourceChangeResult
(Set<MessageQueue> increaseSet, Set<MessageQueue> decreaseSet)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:674
MethodSourceSplitChangeResult
(Set<RocketMQSourceSplit> increaseSet)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:698
MethodTestingEmitterOutput
()
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:70
MethodTimeLagWatermarkGenerator
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/TimeLagWatermarkGenerator.java:32
MethodWaterMarkForAll
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/WaterMarkForAll.java:28
MethodWaterMarkPerQueue
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/WaterMarkPerQueue.java:34
MethodWritableMetadata
( String key, DataType dataType, RocketMQRowDataConverter.Meta
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:300
MethodaddReader
(int subtaskId)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:219
MethodaddSplitsBack
Add a split back to the split enumerator. It will only happen when a {@link SourceReader} fails and there are splits assigned to it after the last suc
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:208
MethodapplyReadableMetadata
(List<String> metadataKeys, DataType producedDataType)
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:158
MethodapplyWritableMetadata
(List<String> metadataKeys, DataType consumedDataType)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:171
MethodasSummaryString
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:188
Methodassign
(Collection<MessageQueue> messageQueues)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:175
Methodassignment
()
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:180
MethodaverageAllocateStrategyTest
()
src/test/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AverageAllocateStrategyTest.java:22
MethodaveragesAllocateStrategyTest
()
src/test/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AverageAllocateStrategyTest.java:43
MethodbroadcastAllocateStrategyTest
()
src/test/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/BroadcastAllocateStrategyTest.java:22
← previousnext →301–400 of 733, ranked by callers