1. 程式人生 > >kafka connect 簡單測試

kafka connect 簡單測試

將檔案中的資料推送topic:connect-test中

配置connect-file-source.properties

cat connect-file-source.properties
name=local-file-source
connector.class=FileStreamSource
tasks.max=1
file=test.txt
topic=connect-test

topic的偏移量儲存在/tmp/connect.offsets這個檔案中,在config/connect-standalone.properties配置,每次connect啟動的時候會根據connector的name獲得topic偏移量,然後在繼續讀取或者寫入資料

檢視檔案test.txt中的資料

cat test.txt

hello 
kafka
hadoop

在使用file source connector時,connector會監聽配置的資料檔案,如果檔案發生變化,例如:追加內容,connector會及時處理新增的資料

啟動kafka connect

bin/connect-standalone.sh config/connect-standalone.properties \
 config/connect-file-source.properties

檢視推送到topic connect-test 的資料

bin/kafka-console
-consumer.sh --bootstrap-server localhost:9092 --topic connect-test --from-beginning {"schema":{"type":"string","optional":false},"payload":"hello "} {"schema":{"type":"string","optional":false},"payload":"kafka"} {"schema":{"type":"string","optional":false},"payload":"hadoop"}

修改connect-file-source.properties,對資料進行轉換,配置如下

name=local-file-source
connector.class=FileStreamSource
tasks.max=1
file=test.txt
topic=connect-test

transforms=MakeMap, InsertSource
transforms.MakeMap.type=org.apache.kafka.connect.transforms.HoistField$Value
transforms.MakeMap.field=line
transforms.InsertSource.type=org.apache.kafka.connect.transforms.InsertField$Value
transforms.InsertSource.static.field=data_source
transforms.InsertSource.static.value=test-file-source

給test.txt增加資料

cat test.txt 
hello 
kafka
hadoop
jafja
spark
hive

啟動kafka connect

bin/connect-standalone.sh config/connect-standalone.properties \
 config/connect-file-source.properties

檢視推送到topic connect-test 的資料

bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic connect-test --from-beginning

{"schema":{"type":"string","optional":false},"payload":"hello "}
{"schema":{"type":"string","optional":false},"payload":"kafka"}
{"schema":{"type":"string","optional":false},"payload":"hadoop"}
{"schema":{"type":"struct","fields":[{"type":"string","optional":false,"field":"line"},{"type":"string","optional":true,"field":"data_source"}],"optional":false},"payload":{"line":"jafja","data_source":"test-file-source"}}
{"schema":{"type":"struct","fields":[{"type":"string","optional":false,"field":"line"},{"type":"string","optional":true,"field":"data_source"}],"optional":false},"payload":{"line":"spark","data_source":"test-file-source"}}
{"schema":{"type":"struct","fields":[{"type":"string","optional":false,"field":"line"},{"type":"string","optional":true,"field":"data_source"}],"optional":false},"payload":{"line":"hive","data_source":"test-file-source"}}

將topic connect-test 的資料儲存到檔案中

配置cat config/connect-file-sink.properties

name=local-file-sink
connector.class=FileStreamSink
tasks.max=1
file=test.sink.txt

啟動kafka connect

bin/connect-standalone.sh config/connect-standalone.properties config/connect-file-sink.properties

檢視生成的test.sink.txt

cat test.sink.txt 
hello 
kafka
hadoop
Struct{line=jafja,data_source=test-file-source}
Struct{line=spark,data_source=test-file-source}
Struct{line=hive,data_source=test-file-source}

啟動多個connector

bin/connect-standalone.sh config/connect-standalone.properties config/connect-file-source.properties config/connect-file-sink.properties