如何配置Confluent Kafka Sink实现数据转换写入Oracle数据库
Got it, let's work through setting up your Confluent Oracle Sink Connector to map your Kafka TEST topic data to the Oracle table exactly as you need. Here's a clear, step-by-step guide:
Before configuring the connector, make sure your target Oracle table exists with the correct columns. Here's a sample DDL you can use (adjust data types based on your actual data):
CREATE TABLE YOUR_SCHEMA.YOUR_TABLE_NAME ( transaction_id VARCHAR2(255) PRIMARY KEY, -- Match your Kafka Key's data type request_object CLOB, -- Use CLOB for JSON/long text; use VARCHAR2 if your requests are short timestamp TIMESTAMP -- Will hold either Kafka's record time or the insertion time );
Replace YOUR_SCHEMA and YOUR_TABLE_NAME with your actual schema and table name. If you're using Oracle 12c+, you can also use the JSON data type for request_object if your payload is JSON.
Let's break down the key settings you'll need to map your Kafka data to the Oracle columns:
Topic & Connection Basics:
topics=TEST: Targets your Kafka topic directly.connection.url: Your Oracle JDBC connection string (e.g.,jdbc:oracle:thin:@//your-host:1521/your-service-name).connection.user/connection.password: Oracle credentials with INSERT/UPDATE permissions on the target table.
Key-to-Column Mapping:
Your Kafka records usetransaction idas the Key. We'll use a Kafka Connect transform to pull this Key into the Value payload so the connector can map it to the Oracletransaction_idcolumn:transforms=extractKey,addTimestamp: Defines the two transforms we'll use.transforms.extractKey.type=org.apache.kafka.connect.transforms.ValueFromKey: Takes the record's Key and adds it to the Value payload.transforms.extractKey.field=transaction_id: Names the new field in the Value that holds the Key (matches your Oracle column name).
Timestamp Handling:
To add the timestamp column, use another transform to inject either the Kafka record's timestamp or the current wall-clock time:transforms.addTimestamp.type=org.apache.kafka.connect.transforms.InsertField$Value: Adds a new field to the Value payload.transforms.addTimestamp.timestamp.field=timestamp: Names the field to match your Oracle column.transforms.addTimestamp.timestamp.type=Record: Uses the Kafka record's creation time. Swap this withWallclockif you want the time the record is inserted into Oracle.
Value & Insert Settings:
key.converter: Useorg.apache.kafka.connect.storage.StringConverterif your Key is a string (adjust if your Key is a numeric type).value.converter: Useorg.apache.kafka.connect.json.JsonConverterif your Value is JSON (withvalue.converter.schemas.enable=falseif you don't use schemas). If you use Avro, switch toio.confluent.connect.avro.AvroConverterand addschema.registry.url.insert.mode=upsert: Handles duplicatetransaction_identries by updating existing rows instead of throwing errors.pk.fields=transaction_id: Marks this column as the primary key for upserts.
Here's a complete connector configuration you can use (replace placeholders with your actual values):
{ "name": "oracle-test-sink-connector", "config": { "connector.class": "io.confluent.connect.oracle.cdc.OracleCdcSinkConnector", "tasks.max": "1", "topics": "TEST", "connection.url": "jdbc:oracle:thin:@//your-oracle-host:1521/your-service-name", "connection.user": "your-oracle-username", "connection.password": "your-oracle-password", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false", "transforms": "extractKey,addTimestamp", "transforms.extractKey.type": "org.apache.kafka.connect.transforms.ValueFromKey", "transforms.extractKey.field": "transaction_id", "transforms.addTimestamp.type": "org.apache.kafka.connect.transforms.InsertField$Value", "transforms.addTimestamp.timestamp.field": "timestamp", "transforms.addTimestamp.timestamp.type": "Record", "table.name.format": "YOUR_SCHEMA.YOUR_TABLE_NAME", "insert.mode": "upsert", "pk.fields": "transaction_id", "auto.create": "false", "auto.evolve": "false" } }
- Data Type Alignment: Ensure your Kafka Key's data type matches the Oracle
transaction_idcolumn (e.g., string ↔ VARCHAR2, number ↔ NUMBER). - Payload Size: If your request objects are large, stick with CLOB or Oracle's JSON type to avoid truncation.
- Permissions: Your Oracle user needs
INSERTandUPDATEpermissions on the target table (since we're using upsert mode). - Schema Handling: If your Value uses a schema (like Avro), update the
value.converterand add the schema registry URL.
内容的提问来源于stack exchange,提问作者user1966395

