MCPcopy Create free account

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

Functions12,620 in github.com/apache/flink-cdc

↓ 2 callersMethodcreateInputSerializer
()
flink-cdc-flink2-compat/src/main/java/org/apache/flink/api/connector/sink2/Sink.java:70
↓ 2 callersMethodcreateJTSLineString
Create JTS line string Geometry from OracleGeometry.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-oracle-cdc/src/main/java/org/apache/flink/cdc/connectors/oracle/source/converters/GeometryConverter.java:247
↓ 2 callersMethodcreateJTSPoint
Create JTS point Geometry from OracleGeometry.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-oracle-cdc/src/main/java/org/apache/flink/cdc/connectors/oracle/source/converters/GeometryConverter.java:220
↓ 2 callersMethodcreateJTSPolygon
Create JTS polygon Geometry from OracleGeometry.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-oracle-cdc/src/main/java/org/apache/flink/cdc/connectors/oracle/source/converters/GeometryConverter.java:346
↓ 2 callersMethodcreateKafkaContainer
This method helps to set commonly used Kafka configurations and aligns the internal Kafka log levels with the ones used by the capturing logger. @par
flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-kafka/src/test/java/org/apache/flink/cdc/connectors/kafka/sink/KafkaUtil.java:59
↓ 2 callersMethodcreateKeyGen
Create a RowDataKeyGen for a table based on its schema. @param schema The table schema @return RowDataKeyGen configured with record key and partition
flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-hudi/src/main/java/org/apache/flink/cdc/connectors/hudi/sink/util/RowDataUtils.java:578
↓ 2 callersMethodcreateMySqlDatabaseSchema
Creates a new {@link MySqlDatabaseSchema} to monitor the latest MySql database schemas.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/debezium/DebeziumUtils.java:106
↓ 2 callersMethodcreateNoStoppingOffset
()
flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/meta/offset/OffsetFactory.java:39
↓ 2 callersMethodcreateNullableExternalConverter
Creates a nullable external converter for the given column type and time zone. @param columnType The type of the column to convert. @param zoneId The
flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-elasticsearch/src/main/java/org/apache/flink/cdc/connectors/elasticsearch/serializer/ElasticsearchRowConverter.java:60
↓ 2 callersMethodcreateNullableExternalConverter
( DataType type, ZoneId pipelineZoneId)
flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-oceanbase/src/main/java/org/apache/flink/cdc/connectors/oceanbase/sink/OceanBaseRowConvert.java:58
↓ 2 callersMethodcreateOceanBaseContainerForJdbc
()
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-oceanbase-cdc/src/test/java/org/apache/flink/cdc/connectors/oceanbase/OceanBaseTestUtils.java:39
↓ 2 callersMethodcreateOrUpdatePublicationModeFilterted
( String tableFilterString, Statement stmt, boolean isUpdate)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql/connection/PostgresReplicationConnection.java:238
↓ 2 callersMethodcreatePartitionMap
( String scheme, String hosts, String database, String collection)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mongodb-cdc/src/main/java/org/apache/flink/cdc/connectors/mongodb/source/utils/MongoRecordUtils.java:179
↓ 2 callersMethodcreatePostgreSqlSource
(int heartbeatInterval)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/test/java/org/apache/flink/cdc/connectors/postgres/PostgreSQLSourceTest.java:629
↓ 2 callersMethodcreatePostgreSqlSourceWithHeartbeatEnabled
()
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/test/java/org/apache/flink/cdc/connectors/postgres/PostgreSQLSourceTest.java:625
↓ 2 callersMethodcreateReader
(SourceReaderContext readerContext)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/apache/flink/cdc/connectors/postgres/source/PostgresSourceBuilder.java:425
↓ 2 callersMethodcreateReader
( SqlServerSourceConfig sourceConfig, SourceReaderContext readerContext, S
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-sqlserver-cdc/src/test/java/org/apache/flink/cdc/connectors/sqlserver/source/reader/SqlserverSourceReaderTest.java:198
↓ 2 callersMethodcreateReader
( OracleSourceConfig configuration)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-oracle-cdc/src/test/java/org/apache/flink/cdc/connectors/oracle/source/reader/OracleSourceReaderTest.java:219
↓ 2 callersMethodcreateRecordEmitterWithCounter
Helper method to create a MySqlRecordEmitter that counts emitted records.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/test/java/org/apache/flink/cdc/connectors/mysql/source/reader/MySqlRecordEmitterTest.java:352
↓ 2 callersMethodcreateReplicationConnection
( PostgresTaskContext taskContext, PostgresConnection postgresConnection,
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql/PostgresObjectUtils.java:117
↓ 2 callersMethodcreateSchemaEventsForTables
( RelationalSnapshotContext<MySqlPartition, MySqlOffsetContext> snapshotContext, final
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/io/debezium/connector/mysql/MySqlSnapshotChangeEventSource.java:462
↓ 2 callersMethodcreateSerializationSchema
In flink>=1.20, the constructor of JsonRowDataSerializationSchema has 6 parameters, and in flink<1.20, the constructor of JsonRowDataSerializationSche
flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-kafka/src/main/java/org/apache/flink/cdc/connectors/kafka/utils/JsonRowDataSerializationSchemaUtils.java:42
↓ 2 callersMethodcreateSourceBuilder
(String... tableNames)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/test/java/org/apache/flink/cdc/connectors/mysql/source/TableExclusionDuringSnapshotIT.java:192
↓ 2 callersMethodcreateSourceFunction
()
flink-cdc-flink2-compat/src/main/java/org/apache/flink/table/connector/source/SourceFunctionProvider.java:66
↓ 2 callersMethodcreateSourceOffsetMap
( final BsonDocument idDocument, boolean isSnapshotRecord)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mongodb-cdc/src/main/java/org/apache/flink/cdc/connectors/mongodb/source/utils/MongoRecordUtils.java:171
↓ 2 callersMethodcreateSourceRecord
( final Map<String, String> partition, final Map<String, String> sourceOffset,
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mongodb-cdc/src/main/java/org/apache/flink/cdc/connectors/mongodb/source/utils/MongoRecordUtils.java:132
↓ 2 callersMethodcreateStarRocksDataSink
(Configuration factoryConfiguration)
flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-starrocks/src/test/java/org/apache/flink/cdc/connectors/starrocks/sink/utils/StarRocksSinkTestBase.java:139
↓ 2 callersMethodcreateStreamSplit
( PostgresSourceConfig sourceConfig, PostgresDialect dialect)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/test/java/org/apache/flink/cdc/connectors/postgres/source/fetch/IncrementalSourceStreamFetcherTest.java:194
↓ 2 callersMethodcreateStreamSplit
()
flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/StreamSplitAssigner.java:170
↓ 2 callersMethodcreateTable
equals to run a sql like: create table table_name (col_name1 type1 comment [, col_name2 type2 ...]);.
flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-maxcompute/src/main/java/org/apache/flink/cdc/connectors/maxcompute/utils/SchemaEvolutionUtils.java:67
↓ 2 callersMethodcreateTableConfig
( Configuration baseConfig, Schema cdcSchema, TableId tableId)
flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-hudi/src/main/java/org/apache/flink/cdc/connectors/hudi/sink/util/ConfigUtils.java:38
↓ 2 callersMethodcreateTableSource
( ResolvedSchema schema, Map<String, String> options)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-tidb-cdc/src/test/java/org/apache/flink/cdc/connectors/tidb/table/TiDBTableSourceFactoryTest.java:137
↓ 2 callersMethodcreateTiKVReadableMetadata
( String database, String tableName)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-tidb-cdc/src/main/java/org/apache/flink/cdc/connectors/tidb/table/TiKVReadableMetadata.java:108
↓ 2 callersMethodcreateWriter
(WriterInitContext writerInitContext)
flink-cdc-flink2-compat/src/main/java/org/apache/flink/api/connector/sink2/Sink.java:41
↓ 2 callersMethodcurrentBsonTimestamp
()
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mongodb-cdc/src/main/java/org/apache/flink/cdc/connectors/mongodb/source/utils/MongoRecordUtils.java:120
↓ 2 callersMethodcurrentMySqlLatestOffset
Gets the latest offset of current MySQL server.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/test/java/org/apache/flink/cdc/connectors/mysql/LegacyMySqlSourceTest.java:1133
↓ 2 callersMethodcurrentOffset
(PostgresConnection jdbcConnection)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql/Utils.java:43
↓ 2 callersMethoddatabase
The name of the PostgreSQL database from which to stream the changes.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/apache/flink/cdc/connectors/postgres/source/PostgresSourceBuilder.java:92
↓ 2 callersMethoddatabase
The name of the SQL Server database from which to stream the changes.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-sqlserver-cdc/src/main/java/org/apache/flink/cdc/connectors/sqlserver/SqlServerSource.java:68
↓ 2 callersMethoddatabase
(String database)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mongodb-cdc/src/main/java/org/apache/flink/cdc/connectors/mongodb/source/offset/ChangeStreamDescriptor.java:80
↓ 2 callersMethoddatabaseExists
(String databaseName)
flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-oceanbase/src/main/java/org/apache/flink/cdc/connectors/oceanbase/catalog/OceanBaseCatalog.java:80
↓ 2 callersMethoddatabaseList
An required list of regular expressions that match database names to be monitored; any database name not included in the whitelist will be excluded fr
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/MySqlSourceBuilder.java:72
↓ 2 callersMethoddatabaseList
A required list of regular expressions that match database names to be monitored; any database name not included in the whitelist will be excluded fro
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-db2-cdc/src/main/java/org/apache/flink/cdc/connectors/db2/source/Db2SourceBuilder.java:67
↓ 2 callersMethoddebeziumProperties
The Debezium Postgres connector properties. For example, "snapshot.mode".
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/apache/flink/cdc/connectors/postgres/source/PostgresSourceBuilder.java:232
↓ 2 callersMethoddebeziumProperties
The Debezium SqlSever connector properties. For example, "snapshot.mode".
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-sqlserver-cdc/src/main/java/org/apache/flink/cdc/connectors/sqlserver/source/SqlServerSourceBuilder.java:186
↓ 2 callersMethoddebeziumProperties
The Debezium MySQL connector properties. For example, "snapshot.mode".
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/MySqlSourceBuilder.java:227
↓ 2 callersMethoddebeziumProperties
The Debezium Db2 connector properties. For example, "snapshot.mode".
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-db2-cdc/src/main/java/org/apache/flink/cdc/connectors/db2/source/Db2SourceBuilder.java:189
↓ 2 callersMethoddebeziumProperties
The Debezium Oracle connector properties. For example, "snapshot.mode".
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-oracle-cdc/src/main/java/org/apache/flink/cdc/connectors/oracle/OracleSource.java:119
↓ 2 callersMethoddebeziumProperties
The Debezium Oracle connector properties. For example, "snapshot.mode".
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-oracle-cdc/src/main/java/org/apache/flink/cdc/connectors/oracle/source/OracleSourceBuilder.java:203
↓ 2 callersMethoddebeziumProperties
The Debezium Vitess connector properties.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-vitess-cdc/src/main/java/org/apache/flink/cdc/connectors/vitess/VitessSource.java:242
↓ 2 callersMethoddecodingPluginName
The name of the Postgres logical decoding plug-in installed on the server. Supported values are decoderbufs, wal2json, wal2json_rds, wal2json_streamin
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/apache/flink/cdc/connectors/postgres/source/PostgresSourceBuilder.java:74
↓ 2 callersMethoddecodingPluginName
The name of the Vitess logical decoding plug-in installed on the server. Supported values are decoderbufs
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-vitess-cdc/src/main/java/org/apache/flink/cdc/connectors/vitess/VitessSource.java:70
↓ 2 callersMethoddeduceMergedCreateTableEvent
Deduce merged CreateTableEvent.
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/common/SchemaDerivator.java:348
↓ 2 callersMethoddeduceSubExpressionType
( List<Column> columns, SqlNode subExpression, List<UserDefinedFunctionDes
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/TransformParser.java:610
↓ 2 callersMethoddefaultEncodeUTF8
(String str, byte[] bytes)
flink-cdc-common/src/main/java/org/apache/flink/cdc/common/utils/StringUtf8Utils.java:114
↓ 2 callersMethoddeployWithApplicationComposer
( PipelineDeploymentExecutor composeExecutor)
flink-cdc-cli/src/main/java/org/apache/flink/cdc/cli/CliExecutor.java:89
↓ 2 callersMethoddeserialize
(byte[] serialized)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/split/FinishedSnapshotSplitInfo.java:141
↓ 2 callersMethoddeserialize
(int version, byte[] serialized)
flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/meta/split/SourceSplitSerializer.java:122
↓ 2 callersMethoddeserialize
Deserialize the TiDB record.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-tidb-cdc/src/main/java/org/apache/flink/cdc/connectors/tidb/TiKVChangeEventDeserializationSchema.java:39
↓ 2 callersMethoddeserialize
(DataInputView source)
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/serializer/schema/DataTypeSerializer.java:188
↓ 2 callersMethoddeserializeBinlogPendingSplitsState
(DataInputDeserializer in)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/assigners/state/PendingSplitsStateSerializer.java:328
↓ 2 callersMethoddeserializeEventIteratorSplit
( DataInputViewStreamWrapper view)
flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/source/ValuesDataSource.java:170
↓ 2 callersMethoddeserializeLegacySnapshotPendingSplitsState
( int splitVersion, DataInputDeserializer in)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/assigners/state/PendingSplitsStateSerializer.java:196
↓ 2 callersMethoddeserializeLegacySnapshotPendingSplitsState
( int splitVersion, DataInputDeserializer in)
flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/state/PendingSplitsStateSerializer.java:213
↓ 2 callersMethoddeserializeMessages
( ByteBuffer buffer, ReplicationMessageProcessor processor)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql/connection/PostgresReplicationConnection.java:616
↓ 2 callersMethoddeserializeReuse
(BinaryMapData reuse, DataInputView source)
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/serializer/data/MapDataSerializer.java:198
↓ 2 callersMethoddeserializeSnapshotPendingSplitsState
( int version, int splitVersion, DataInputDeserializer in)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/assigners/state/PendingSplitsStateSerializer.java:250
↓ 2 callersMethoddeserializeSnapshotPendingSplitsState
( int version, int splitVersion, DataInputDeserializer in)
flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/state/PendingSplitsStateSerializer.java:267
↓ 2 callersMethoddeserializeStreamPendingSplitsState
(DataInputDeserializer in)
flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/state/PendingSplitsStateSerializer.java:350
↓ 2 callersMethoddeserializer
The deserializer used to convert from consumed {@link org.apache.kafka.connect.source.SourceRecord}.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/apache/flink/cdc/connectors/postgres/source/PostgresSourceBuilder.java:258
↓ 2 callersMethoddeserializer
The deserializer used to convert from consumed {@link org.apache.kafka.connect.source.SourceRecord}.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-sqlserver-cdc/src/main/java/org/apache/flink/cdc/connectors/sqlserver/SqlServerSource.java:107
↓ 2 callersMethoddeserializer
The deserializer used to convert from consumed {@link org.apache.kafka.connect.source.SourceRecord}.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-sqlserver-cdc/src/main/java/org/apache/flink/cdc/connectors/sqlserver/source/SqlServerSourceBuilder.java:195
↓ 2 callersMethoddeserializer
The deserializer used to convert from consumed {@link org.apache.kafka.connect.source.SourceRecord}.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/MySqlSourceBuilder.java:236
↓ 2 callersMethoddeserializer
(DebeziumDeserializationSchema<T> deserializer)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-db2-cdc/src/main/java/org/apache/flink/cdc/connectors/db2/Db2Source.java:125
↓ 2 callersMethoddeserializer
The deserializer used to convert from consumed {@link org.apache.kafka.connect.source.SourceRecord}.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-db2-cdc/src/main/java/org/apache/flink/cdc/connectors/db2/source/Db2SourceBuilder.java:198
↓ 2 callersMethoddeserializer
The deserializer used to convert from consumed {@link org.apache.kafka.connect.source.SourceRecord}.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-oracle-cdc/src/main/java/org/apache/flink/cdc/connectors/oracle/source/OracleSourceBuilder.java:227
↓ 2 callersMethoddeserializer
The deserializer used to convert from consumed {@link org.apache.kafka.connect.source.SourceRecord}.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-vitess-cdc/src/main/java/org/apache/flink/cdc/connectors/vitess/VitessSource.java:251
↓ 2 callersMethoddeserializer
The deserializer used to convert from consumed {@link org.apache.kafka.connect.source.SourceRecord}.
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mongodb-cdc/src/main/java/org/apache/flink/cdc/connectors/mongodb/source/MongoDBSourceBuilder.java:242
↓ 2 callersMethoddestroyHarness
()
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/transform/TransformOperatorWithSchemaEvolveTest.java:218
↓ 2 callersMethoddiscoverDataCollectionSchemas
Discovers the captured data collections' schema by {@link SourceConfig}. @param sourceConfig a basic source configuration.
flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/dialect/DataSourceDialect.java:56
↓ 2 callersMethoddoReadTableColumn
( ResultSet columnMetadata, TableId tableId, Tables.ColumnNameFilter columnFilter)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql/connection/PostgresConnection.java:703
↓ 2 callersMethoddrainSnapshotSplits
(MySqlHybridSplitAssigner assigner)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/test/java/org/apache/flink/cdc/connectors/mysql/source/assigners/MySqlHybridSplitAssignerTest.java:265
↓ 2 callersMethoddropDorisDatabase
(String databaseName)
flink-cdc-e2e-tests/flink-cdc-pipeline-e2e-tests/src/test/java/org/apache/flink/cdc/pipeline/tests/MySqlToDorisE2eITCase.java:966
↓ 2 callersMethoddropReplicationSlot
Drops a replication slot that was created on the DB @param slotName the name of the replication slot, may not be null @return {@code true} if the slo
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql/connection/PostgresConnection.java:486
↓ 2 callersMethodduplicate
()
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/serializer/MapSerializer.java:87
↓ 2 callersMethodemitElement
(SourceRecord element, SourceOutput<SourceRecord> output)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/test/java/org/apache/flink/cdc/connectors/mysql/source/reader/MySqlSourceReaderTest.java:812
↓ 2 callersMethodemitElement
(SourceRecord element, SourceOutput<T> output)
flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/reader/IncrementalSourceRecordEmitter.java:159
↓ 2 callersMethodemitLatestSchema
(TableId tableId)
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/sink/DataSinkOperatorAdapter.java:157
↓ 2 callersMethodemitLatestSchema
(TableId tableId)
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/sink/DataSinkFunctionOperator.java:141
↓ 2 callersMethodemitLatestSchema
(TableId tableId)
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/sink/DataSinkWriterOperator.java:234
↓ 2 callersMethodenableIgnoreNullFields
flink>=1.20 only has the ENCODE_IGNORE_NULL_FIELDS parameter.
flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-kafka/src/main/java/org/apache/flink/cdc/connectors/kafka/utils/JsonRowDataSerializationSchemaUtils.java:134
↓ 2 callersMethodencodeUTF8
This method must have the same result with JDK's String.getBytes.
flink-cdc-common/src/main/java/org/apache/flink/cdc/common/utils/StringUtf8Utils.java:40
↓ 2 callersMethodencodeValue
(String value)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mongodb-cdc/src/main/java/org/apache/flink/cdc/connectors/mongodb/internal/MongoDBEnvelope.java:149
↓ 2 callersMethodendIndex
Get the index in the raw stream past the last character in the token. @return the ending index of the token, which is past the last character
flink-cdc-common/src/main/java/org/apache/flink/cdc/common/text/TokenStream.java:523
↓ 2 callersMethodensureJmLeaderServiceExists
( HaLeadershipControl leadershipControl, JobID jobId)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/test/java/org/apache/flink/cdc/connectors/postgres/testutils/PostgresTestUtils.java:75
↓ 2 callersMethodensureJmLeaderServiceExists
( HaLeadershipControl leadershipControl, JobID jobId)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/test/java/org/apache/flink/cdc/connectors/mysql/source/MySqlSourceTestBase.java:153
↓ 2 callersMethodensureJmLeaderServiceExists
( HaLeadershipControl leadershipControl, JobID jobId)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-oracle-cdc/src/test/java/org/apache/flink/cdc/connectors/oracle/testutils/OracleTestUtils.java:65
↓ 2 callersMethodensureJmLeaderServiceExists
( HaLeadershipControl leadershipControl, JobID jobId)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mongodb-cdc/src/test/java/org/apache/flink/cdc/connectors/mongodb/utils/MongoDBTestUtils.java:182
↓ 2 callersMethodensurePkNonNull
(Schema schema)
flink-cdc-common/src/main/java/org/apache/flink/cdc/common/utils/SchemaUtils.java:408
↓ 2 callersMethodequalTo
(VersionComparable version)
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/utils/VersionComparable.java:77
← previousnext →2,601–2,700 of 12,620, ranked by callers