MCPcopy Create free account

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

Functions733 in github.com/apache/rocketmq-flink

↓ 1 callersMethodgetCatalogTableForSchema
( String topic, GetSchemaResponse getSchemaResponse)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:200
↓ 1 callersMethodgetConsumerGroup
Get the consumer group of the consumer.
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:36
↓ 1 callersMethodgetConsumerProps
Source Config @return properties
src/test/java/org/apache/flink/connector/rocketmq/example/LegacyConnectorExample.java:49
↓ 1 callersMethodgetConsumerProps
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:219
↓ 1 callersMethodgetCurrentWatermark
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/WaterMarkForAll.java:38
↓ 1 callersMethodgetDatabase
(String databaseName)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:142
↓ 1 callersMethodgetDecreaseSet
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:683
↓ 1 callersMethodgetEventTime
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 callersMethodgetFactory
()
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:103
↓ 1 callersMethodgetFunction
(ObjectPath functionPath)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:337
↓ 1 callersMethodgetMailboxExecutor
()
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/context/RocketMQSinkContext.java:53
↓ 1 callersMethodgetPartition
(ObjectPath tablePath, CatalogPartitionSpec partitionSpec)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:281
↓ 1 callersMethodgetPartitionColumnStatistics
( ObjectPath tablePath, CatalogPartitionSpec partitionSpec)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:448
↓ 1 callersMethodgetPartitionStatistics
( ObjectPath tablePath, CatalogPartitionSpec partitionSpec)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:441
↓ 1 callersMethodgetProducerProps
Sink Config @return properties
src/test/java/org/apache/flink/connector/rocketmq/example/LegacyConnectorExample.java:69
↓ 1 callersMethodgetProducerProps
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:238
↓ 1 callersMethodgetSourceChangeResult
(Set<MessageQueue> latestSet)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:718
↓ 1 callersMethodgetSourceSplit
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 callersMethodgetSplitId
(MessageQueue mq)
src/main/java/org/apache/flink/connector/rocketmq/source/util/UtilAll.java:28
↓ 1 callersMethodgetSplitOwner
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 callersMethodgetSplitOwner
(String topic, int partition, int numReaders)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AverageAllocateStrategy.java:35
↓ 1 callersMethodgetSplitReader
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceFetcherManager.java:99
↓ 1 callersMethodgetStoreSize
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 callersMethodgetStrategy
( Configuration rocketmqSourceOptions, SplitEnumeratorContext<RocketMQSourceSplit> con
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/AllocateStrategyFactory.java:36
↓ 1 callersMethodgetTable
(ObjectPath tablePath)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:178
↓ 1 callersMethodgetTableColumnStatistics
(ObjectPath tablePath)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:435
↓ 1 callersMethodgetTableStatistics
(ObjectPath tablePath)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:429
↓ 1 callersMethodgetTransactionProducer
Lazy initialize this backend transaction client.
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/RocketMQCommitter.java:86
↓ 1 callersMethodgetValue
(String[] data, String line, int index)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:174
↓ 1 callersMethodgetVersion
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumStateSerializer.java:45
↓ 1 callersMethodhandleException
(GenericRowData row, int index, Object[] data, Exception e)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:197
↓ 1 callersMethodhandleException
(GenericRowData row, int index, Object[] data, Exception e)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:280
↓ 1 callersMethodhandleFieldIncrement
(String[] data)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:259
↓ 1 callersMethodhandleFieldIncrement
(String[] data)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:357
↓ 1 callersMethodhandleFieldMissing
(String[] data)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:234
↓ 1 callersMethodhandleFieldMissing
(String[] data)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:317
↓ 1 callersMethodhandleInitAssignEvent
(int taskId, SourceInitAssignEvent initAssignEvent)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:588
↓ 1 callersMethodhandleOffsetEvent
(int taskId, SourceReportOffsetEvent sourceReportOffsetEvent)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:620
↓ 1 callersMethodhandleSourceCheckEvent
(SourceCheckEvent sourceEvent)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:198
↓ 1 callersMethodhandleSourceDetectEvent
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:176
↓ 1 callersMethodinitializeSourceSplits
(SourceChangeResult sourceChangeResult)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:319
↓ 1 callersMethodisAllHeaderField
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:194
↓ 1 callersMethodisBounded
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:208
↓ 1 callersMethodisByteArrayType
(String fieldName)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:186
↓ 1 callersMethodisByteArrayType
(String fieldName)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:269
↓ 1 callersMethodisOnlyHaveVarbinaryDataField
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:125
↓ 1 callersMethodisOnlyHaveVarbinaryDataField
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:186
↓ 1 callersMethodlistDatabases
()
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:137
↓ 1 callersMethodlistFunctions
(String dbName)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:331
↓ 1 callersMethodlistPartitionsByFilter
( ObjectPath tablePath, List<Expression> expressions)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:393
↓ 1 callersMethodlistTables
(String databaseName)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:160
↓ 1 callersMethodlistViews
(String databaseName)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:374
↓ 1 callersMethodnotifyCheckpointComplete
(Map<MessageQueue, Long> offsetsToCommit)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:238
↓ 1 callersMethodoffsetsForTimes
List max offsets for the specified MessageQueues.
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelector.java:78
↓ 1 callersMethodopen
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 callersMethodoptionalOptions
()
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactory.java:60
↓ 1 callersMethodparseAndSetRequiredProperties
()
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:169
↓ 1 callersMethodparseAndSetRequiredProperties
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:200
↓ 1 callersMethodpartitionExists
(ObjectPath tablePath, CatalogPartitionSpec partitionSpec)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:287
↓ 1 callersMethodprepareForRead
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:301
↓ 1 callersMethodprepareSourceEnumeratorState
()
src/test/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumStateSerializerTest.java:45
↓ 1 callersMethodrecordsForSplit
return records container
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:293
↓ 1 callersMethodregisterNewMessageQueue
(MessageQueue messageQueue)
src/main/java/org/apache/flink/connector/rocketmq/source/metrics/RocketMQSourceReaderMetrics.java:47
↓ 1 callersMethodregisterOutBps
(RuntimeContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:62
↓ 1 callersMethodregisterOutLatency
(RuntimeContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:72
↓ 1 callersMethodregisterOutTps
(RuntimeContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:52
↓ 1 callersMethodregisterSinkInTps
(RuntimeContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:42
↓ 1 callersMethodregisterSplits
(RocketMQSourceSplit split)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:250
↓ 1 callersMethodrenameTable
(ObjectPath tablePath, String newTableName, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:387
↓ 1 callersMethodrequiredOptions
()
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactory.java:55
↓ 1 callersMethodresume
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 callersMethodrun
(SourceContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:250
↓ 1 callersMethodsanityCheck
()
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:167
↓ 1 callersMethodsanityCheck
()
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:198
↓ 1 callersMethodseekCommittedOffset
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 callersMethodseekMaxOffset
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 callersMethodseekMinOffset
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 callersMethodseekOffsetByTimestamp
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 callersMethodselect
(List<MessageQueue> mqs, Message msg, Object arg)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/HashMessageQueueSelector.java:27
↓ 1 callersMethodselect
(List<MessageQueue> mqs, Message msg, Object arg)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/RandomMessageQueueSelector.java:30
↓ 1 callersMethodsendMessageInTransaction
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 callersMethodserialize
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 callersMethodserialize
(RocketMQSourceEnumState enumState)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumStateSerializer.java:50
↓ 1 callersMethodserializeKey
(Map tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleKeyValueSerializationSchema.java:44
↓ 1 callersMethodserializeValue
(Map tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleKeyValueSerializationSchema.java:53
↓ 1 callersMethodsetAssignedMq
(Map<MessageQueue, Tuple2<Long, Long>> assignedMq)
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceCheckEvent.java:35
↓ 1 callersMethodsetBroker
(String broker)
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:29
↓ 1 callersMethodsetBrokerName
(String brokerName)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:69
↓ 1 callersMethodsetCheckpoint
(long checkpoint)
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:33
↓ 1 callersMethodsetColumnErrorDebug
(boolean columnErrorDebug)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:329
↓ 1 callersMethodsetColumnErrorDebug
(boolean columnErrorDebug)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:505
↓ 1 callersMethodsetCurrentOffset
(long currentOffset)
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQSourceSplitState.java:40
↓ 1 callersMethodsetData
(byte[] data)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/BytesMessage.java:35
↓ 1 callersMethodsetDeliveryGuarantee
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 callersMethodsetDeserializer
( RocketMQDeserializationSchema<OUT> recordDeserializer)
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:128
↓ 1 callersMethodsetEncoding
(String encoding)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:319
↓ 1 callersMethodsetEncoding
(String encoding)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:490
↓ 1 callersMethodsetFieldDelimiter
(String fieldDelimiter)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:324
↓ 1 callersMethodsetFieldDelimiter
(String fieldDelimiter)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:495
↓ 1 callersMethodsetHasMetadata
(boolean hasMetadata)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:510
← previousnext →201–300 of 733, ranked by callers