Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/apache/rocketmq-flink
/ functions
Functions
733 in github.com/apache/rocketmq-flink
⨍
Functions
733
◇
Types & classes
161
↓ 1 callers
Method
getCatalogTableForSchema
( String topic, GetSchemaResponse getSchemaResponse)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:200
↓ 1 callers
Method
getConsumerGroup
Get the consumer group of the consumer.
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:36
↓ 1 callers
Method
getConsumerProps
Source Config @return properties
src/test/java/org/apache/flink/connector/rocketmq/example/LegacyConnectorExample.java:49
↓ 1 callers
Method
getConsumerProps
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:219
↓ 1 callers
Method
getCurrentWatermark
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/WaterMarkForAll.java:38
↓ 1 callers
Method
getDatabase
(String databaseName)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:142
↓ 1 callers
Method
getDecreaseSet
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:683
↓ 1 callers
Method
getEventTime
Get the event time of the message, which is used for filtering and sorting. @return the event time
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:101
↓ 1 callers
Method
getFactory
()
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:103
↓ 1 callers
Method
getFunction
(ObjectPath functionPath)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:337
↓ 1 callers
Method
getMailboxExecutor
()
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContext.java:53
↓ 1 callers
Method
getPartition
(ObjectPath tablePath, CatalogPartitionSpec partitionSpec)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:281
↓ 1 callers
Method
getPartitionColumnStatistics
( ObjectPath tablePath, CatalogPartitionSpec partitionSpec)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:448
↓ 1 callers
Method
getPartitionStatistics
( ObjectPath tablePath, CatalogPartitionSpec partitionSpec)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:441
↓ 1 callers
Method
getProducerProps
Sink Config @return properties
src/test/java/org/apache/flink/connector/rocketmq/example/LegacyConnectorExample.java:69
↓ 1 callers
Method
getProducerProps
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:238
↓ 1 callers
Method
getSourceChangeResult
(Set<MessageQueue> latestSet)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:718
↓ 1 callers
Method
getSourceSplit
Use the current offset as the starting offset to create a new RocketMQSourceSplit. @return a new RocketMQSourceSplit which uses the current offset as
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplitState.java:49
↓ 1 callers
Method
getSplitId
(MessageQueue mq)
src/main/java/org/apache/flink/connector/rocketmq/source/util/UtilAll.java:28
↓ 1 callers
Method
getSplitOwner
Returns the index of the target subtask that a specific queue should be assigned to.
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/ConsistentHashAllocateStrategy.java:35
↓ 1 callers
Method
getSplitOwner
(String topic, int partition, int numReaders)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AverageAllocateStrategy.java:35
↓ 1 callers
Method
getSplitReader
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceFetcherManager.java:99
↓ 1 callers
Method
getStoreSize
Get the size of the message in bytes. @return the message size
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:80
↓ 1 callers
Method
getStrategy
( Configuration rocketmqSourceOptions, SplitEnumeratorContext<RocketMQSourceSplit> con
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AllocateStrategyFactory.java:36
↓ 1 callers
Method
getTable
(ObjectPath tablePath)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:178
↓ 1 callers
Method
getTableColumnStatistics
(ObjectPath tablePath)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:435
↓ 1 callers
Method
getTableStatistics
(ObjectPath tablePath)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:429
↓ 1 callers
Method
getTransactionProducer
Lazy initialize this backend transaction client.
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/RocketMQCommitter.java:86
↓ 1 callers
Method
getValue
(String[] data, String line, int index)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:174
↓ 1 callers
Method
getVersion
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumStateSerializer.java:45
↓ 1 callers
Method
handleException
(GenericRowData row, int index, Object[] data, Exception e)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:197
↓ 1 callers
Method
handleException
(GenericRowData row, int index, Object[] data, Exception e)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:280
↓ 1 callers
Method
handleFieldIncrement
(String[] data)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:259
↓ 1 callers
Method
handleFieldIncrement
(String[] data)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:357
↓ 1 callers
Method
handleFieldMissing
(String[] data)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:234
↓ 1 callers
Method
handleFieldMissing
(String[] data)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:317
↓ 1 callers
Method
handleInitAssignEvent
(int taskId, SourceInitAssignEvent initAssignEvent)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:588
↓ 1 callers
Method
handleOffsetEvent
(int taskId, SourceReportOffsetEvent sourceReportOffsetEvent)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:620
↓ 1 callers
Method
handleSourceCheckEvent
(SourceCheckEvent sourceEvent)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:198
↓ 1 callers
Method
handleSourceDetectEvent
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:176
↓ 1 callers
Method
initializeSourceSplits
(SourceChangeResult sourceChangeResult)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:319
↓ 1 callers
Method
isAllHeaderField
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:194
↓ 1 callers
Method
isBounded
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:208
↓ 1 callers
Method
isByteArrayType
(String fieldName)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:186
↓ 1 callers
Method
isByteArrayType
(String fieldName)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:269
↓ 1 callers
Method
isOnlyHaveVarbinaryDataField
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:125
↓ 1 callers
Method
isOnlyHaveVarbinaryDataField
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:186
↓ 1 callers
Method
listDatabases
()
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:137
↓ 1 callers
Method
listFunctions
(String dbName)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:331
↓ 1 callers
Method
listPartitionsByFilter
( ObjectPath tablePath, List<Expression> expressions)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:393
↓ 1 callers
Method
listTables
(String databaseName)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:160
↓ 1 callers
Method
listViews
(String databaseName)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:374
↓ 1 callers
Method
notifyCheckpointComplete
(Map<MessageQueue, Long> offsetsToCommit)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:238
↓ 1 callers
Method
offsetsForTimes
List max offsets for the specified MessageQueues.
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelector.java:78
↓ 1 callers
Method
open
Initialization method for the schema. It is called before the actual working methods {@link #deserialize} and thus suitable for one time setup work.
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchema.java:44
↓ 1 callers
Method
optionalOptions
()
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactory.java:60
↓ 1 callers
Method
parseAndSetRequiredProperties
()
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:169
↓ 1 callers
Method
parseAndSetRequiredProperties
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:200
↓ 1 callers
Method
partitionExists
(ObjectPath tablePath, CatalogPartitionSpec partitionSpec)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:287
↓ 1 callers
Method
prepareForRead
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:301
↓ 1 callers
Method
prepareSourceEnumeratorState
()
src/test/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumStateSerializerTest.java:45
↓ 1 callers
Method
recordsForSplit
return records container
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:293
↓ 1 callers
Method
registerNewMessageQueue
(MessageQueue messageQueue)
src/main/java/org/apache/flink/connector/rocketmq/source/metrics/RocketMQSourceReaderMetrics.java:47
↓ 1 callers
Method
registerOutBps
(RuntimeContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:62
↓ 1 callers
Method
registerOutLatency
(RuntimeContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:72
↓ 1 callers
Method
registerOutTps
(RuntimeContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:52
↓ 1 callers
Method
registerSinkInTps
(RuntimeContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:42
↓ 1 callers
Method
registerSplits
(RocketMQSourceSplit split)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:250
↓ 1 callers
Method
renameTable
(ObjectPath tablePath, String newTableName, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:387
↓ 1 callers
Method
requiredOptions
()
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactory.java:55
↓ 1 callers
Method
resume
Resuming message pulling from the message queues. @param messageQueues message queues that need to be resumed.
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:85
↓ 1 callers
Method
run
(SourceContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:250
↓ 1 callers
Method
sanityCheck
()
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:167
↓ 1 callers
Method
sanityCheck
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:198
↓ 1 callers
Method
seekCommittedOffset
Seek consumer group previously committed offset @param messageQueue rocketmq queue to locate single queue @return offset for message queue
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:103
↓ 1 callers
Method
seekMaxOffset
Seek consumer group previously committed offset @param messageQueue rocketmq queue to locate single queue @return offset for message queue
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:119
↓ 1 callers
Method
seekMinOffset
Seek consumer group previously committed offset @param messageQueue rocketmq queue to locate single queue @return offset for message queue
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:111
↓ 1 callers
Method
seekOffsetByTimestamp
Seek consumer group previously committed offset @param messageQueue rocketmq queue to locate single queue @return offset for message queue
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:127
↓ 1 callers
Method
select
(List<MessageQueue> mqs, Message msg, Object arg)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/HashMessageQueueSelector.java:27
↓ 1 callers
Method
select
(List<MessageQueue> mqs, Message msg, Object arg)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/RandomMessageQueueSelector.java:30
↓ 1 callers
Method
sendMessageInTransaction
Sends the message to the messaging system and returns a Future for the send operation. @param message the message to be sent @return a Future for the
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducer.java:63
↓ 1 callers
Method
serialize
Serializes given element and returns it as a {@link MessageExt}. @param element element to be serialized @param context context to possibly determine
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/serializer/RocketMQSerializationSchema.java:43
↓ 1 callers
Method
serialize
(RocketMQSourceEnumState enumState)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumStateSerializer.java:50
↓ 1 callers
Method
serializeKey
(Map tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleKeyValueSerializationSchema.java:44
↓ 1 callers
Method
serializeValue
(Map tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleKeyValueSerializationSchema.java:53
↓ 1 callers
Method
setAssignedMq
(Map<MessageQueue, Tuple2<Long, Long>> assignedMq)
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceCheckEvent.java:35
↓ 1 callers
Method
setBroker
(String broker)
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:29
↓ 1 callers
Method
setBrokerName
(String brokerName)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:69
↓ 1 callers
Method
setCheckpoint
(long checkpoint)
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:33
↓ 1 callers
Method
setColumnErrorDebug
(boolean columnErrorDebug)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:329
↓ 1 callers
Method
setColumnErrorDebug
(boolean columnErrorDebug)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:505
↓ 1 callers
Method
setCurrentOffset
(long currentOffset)
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplitState.java:40
↓ 1 callers
Method
setData
(byte[] data)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/BytesMessage.java:35
↓ 1 callers
Method
setDeliveryGuarantee
Sets the wanted the {@link DeliveryGuarantee}. @param deliveryGuarantee target delivery guarantee @return {@link RocketMQSinkBuilder}
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:89
↓ 1 callers
Method
setDeserializer
( RocketMQDeserializationSchema<OUT> recordDeserializer)
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:128
↓ 1 callers
Method
setEncoding
(String encoding)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:319
↓ 1 callers
Method
setEncoding
(String encoding)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:490
↓ 1 callers
Method
setFieldDelimiter
(String fieldDelimiter)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:324
↓ 1 callers
Method
setFieldDelimiter
(String fieldDelimiter)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:495
↓ 1 callers
Method
setHasMetadata
(boolean hasMetadata)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:510
← previous
next →
201–300 of 733, ranked by callers