Using Kafka Connect
1. Understanding Connect Framework
| Concept | Description | Detail |
|---|---|---|
| Connect | Scalable data integration | No code |
| Source | External → Kafka | Ingest |
| Sink | Kafka → External | Export |
2. Running Standalone Mode
| Aspect | Description | Detail |
|---|---|---|
| Single Process | One worker | Dev/testing |
| Offsets | Local file | Not fault-tolerant |
| Config | Properties files | CLI args |
Example: Start standalone
connect-standalone.sh config/connect-standalone.properties \
config/file-source.properties
3. Running Distributed Mode
| Aspect | Description | Detail |
|---|---|---|
| Multi-Worker | Cluster of workers | Scalable |
| Offsets | Kafka topics | Fault-tolerant |
| Config | REST API | No file connectors |
4. Creating Source Connectors
| Field | Description | Detail |
|---|---|---|
| connector.class | Source implementation | e.g. JDBC |
| tasks.max | Parallel tasks | Throughput |
| topic.prefix | Output topic naming | Per source |
Example: JDBC source
{"name":"jdbc-src","config":{
"connector.class":"io.confluent.connect.jdbc.JdbcSourceConnector",
"connection.url":"jdbc:postgresql://db/app","mode":"incrementing"}}
5. Creating Sink Connectors
| Field | Description | Detail |
|---|---|---|
| connector.class | Sink implementation | e.g. S3 |
| topics | Input topics | Comma list |
| flush.size | Batch before write | Throughput |
Example: S3 sink
{"name":"s3-sink","config":{
"connector.class":"io.confluent.connect.s3.S3SinkConnector",
"topics":"orders","s3.bucket.name":"my-bucket","flush.size":"1000"}}
6. Configuring Connector Properties
| Property | Description | Detail |
|---|---|---|
| key.converter | Key format | Override worker |
| transforms | SMT chain | Inline transform |
| errors.tolerance | none / all | DLQ routing |
Example: SMT transform
{"transforms":"route",
"transforms.route.type":"org.apache.kafka.connect.transforms.RegexRouter",
"transforms.route.regex":"(.*)","transforms.route.replacement":"prefix_$1"}
7. Managing Connectors via REST API
| Endpoint | Method | Action |
|---|---|---|
| /connectors | POST | Create |
| /connectors/{n}/config | PUT | Update |
| /connectors/{n}/restart | POST | Restart |
Example: Create via REST
curl -X POST -H "Content-Type: application/json" \
--data @connector.json http://localhost:8083/connectors
8. Listing Connectors
| Endpoint | Description | Detail |
|---|---|---|
| GET /connectors | All connector names | JSON array |
| ?expand=status | Include status | Detailed |
| ?expand=info | Include config | Detailed |
9. Checking Connector Status
| State | Description | Detail |
|---|---|---|
| RUNNING | Healthy | Normal |
| FAILED | Error with trace | Inspect |
| PAUSED | Manually stopped | Resumable |
10. Pausing and Resuming Connectors
| Endpoint | Method | Action |
|---|---|---|
| /{n}/pause | PUT | Stop tasks |
| /{n}/resume | PUT | Restart tasks |
| /{n}/stop | PUT | Stop + free NEW |
11. Deleting Connectors
| Aspect | Description | Detail |
|---|---|---|
| DELETE /{n} | Remove connector | Stops tasks |
| Offsets | Retained in topics | Resume later |
| Idempotent | 404 if absent | Safe |