Kafka是一个分布式的流处理平台

Kafka提供封装好的客户端以方便开发者连接服务器,目前常用的客户端有两种:
1<dependency> 2 <groupId>org.apache.kafka</groupId> 3 <artifactId>kafka-clients</artifactId> 4 <version>0.10.0.1</version> 5</dependency> 6 7 8<dependency> 9 <groupId>org.apache.kafka</groupId> 10 <artifactId>kafka_2.10</artifactId> 11 <version>0.10.0.1</version> 12</dependency>
上面的一种为官方目前推荐的(但貌似很多生产环境是用下边一种的),具体的历史原因第二种2013年就已经问世了,第一种是在2015年的时候才出现的,为了响应官方号召我将主要分析第一种客户端
1. Producers
Producer API 允许应用程序发布一个或多个流式数据给topic, 如果需要使用这个API,你需要在maven的pom配置文件中,增加一个依赖
1 <dependency> 2 <groupId>org.apache.kafka</groupId> 3 <artifactId>kafka-clients</artifactId> 4 <version>0.10.1.0</version> 5 </dependency>
我们会看到客户端(刚才添加这个东西会引入一个kafka客户端)中定义了一个Producer接口,并继承于Closeable接口(这个接口源于JDK1.5,通常用于需要关闭对象的资源,其中有一个close方法用于关闭,并抛出一个IOException),我们简单看一下类图,Apache已经为我们搞定了两个实现,如下:

KafkaProducer
这个是Apache提供的一个Producer接口实现,同时标明了他是单例且线程安全的,使用方法大致如下:
1 Properties props = new Properties(); 2 props.put("bootstrap.servers", "localhost:9092"); 3 props.put("acks", "all"); 4 props.put("retries", 0); 5 props.put("batch.size", 16384); 6 props.put("linger.ms", 1); 7 props.put("buffer.memory", 33554432); 8 props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); 9 props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); 10 11 Producer<String, String> producer = new KafkaProducer(props); 12 for (int i = 0; i < 100; i++) 13 producer.send(new ProducerRecord<String, String>("my-topic", Integer.toString(i), Integer.toString(i))); 14 15 producer.close();
Properties的配置参数大家可以到 http://kafka.apache.org/documentation.html#producerconfigs 查看手册
MockProducer
这个是apache为开发测试提供的一个实现类,他提供了一些扩展的方便,可以方便我们开发的时候用于调试
2. Consumers
消费者的接口设计稍微复杂一些,应该是为了方便使用,同样apache也提供了两个实现,如果需要使用apache提供的客户端,同样需要引入如下包
1 <dependency> 2 <groupId>org.apache.kafka</groupId> 3 <artifactId>kafka-clients</artifactId> 4 <version>0.10.1.0</version> 5 </dependency>
Consumer接口看以下类图:

KafkaConsumer
需要注意的一点是这个KafkaConsumer实现,并不是线程安全的,大家在使用过程中,需要保证线程安全,
这个Consumer的实现,我们可以参考如下几种用法:
1 Properties props = new Properties(); 2 props.put("bootstrap.servers", "localhost:9092"); 3 props.put("group.id", "test"); 4 props.put("enable.auto.commit", "true"); 5 props.put("auto.commit.interval.ms", "1000"); 6 props.put("session.timeout.ms", "30000"); 7 props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); 8 props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); 9 KafkaConsumer<String, String> consumer = new KafkaConsumer(props); 10 consumer.subscribe(Arrays.asList("foo", "bar")); 11 while(true){ 12 ConsumerRecords<String, String> records = consumer.poll(100); 13 for (ConsumerRecord<String, String> record : records) 14 System.out.printf("offset =%d, key =%s, value =%s", record.offset(), record.key(), record.value()); 15 }
上面这种用法设置了 enable.auto.commit = true, offsets 将会根据 auto.commit.interval.ms 的值,进行自动提交
1 Properties props = new Properties(); 2 props.put("bootstrap.servers", "localhost:9092"); 3 props.put("group.id", "test"); 4 props.put("enable.auto.commit", "false"); 5 props.put("auto.commit.interval.ms", "1000"); 6 props.put("session.timeout.ms", "30000"); 7 props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); 8 props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); 9 KafkaConsumer<String, String> consumer = new KafkaConsumer(props); 10 consumer.subscribe(Arrays.asList("foo", "bar")); 11 final int minBatchSize = 200; 12 List<ConsumerRecord<String, String>> buffer = new ArrayList<ConsumerRecord<String, String>>(); 13 while (true) { 14 ConsumerRecords<String, String> records = consumer.poll(100); 15 for (ConsumerRecord<String, String> record : records) { 16 buffer.add(record); 17 } 18 if (buffer.size() >= minBatchSize) { 19 insertIntoDb(buffer); //这里插入数据库 20 consumer.commitSync(); 21 buffer.clear(); 22 } 23 }
以上这种用法关闭了 enable.auto.commit 参数,通过 consumer的commitSync方法来手动commit(Commit操作表示客户端通知服务端,信息已经收到,服务器会标记信息为 commited 状态),同时使用List做了一个缓存,来批量进行数据库写入操作
1public class KafkaConsumerRunner implements Runnable { 2 3 public static void main(String[] args) throws InterruptedException { 4 5 Properties props = new Properties(); 6 props.put("bootstrap.servers", "192.168.3.37:9092"); 7 props.put("group.id", "test"); 8 props.put("enable.auto.commit", "false"); 9 props.put("auto.commit.interval.ms", "1000"); 10 props.put("session.timeout.ms", "30000"); 11 props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); 12 props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); 13 14 int totleThread = 5; 15 CountDownLatch countDown = new CountDownLatch(totleThread); 16 17 List<KafkaConsumerRunner> consumers = startServer(props, totleThread, countDown); 18 19 Thread.sleep(10000); 20 21 stopServer(consumers); 22 23 countDown.await(); 24 System.out.println("shutdown all done"); 25 System.exit(0); 26 } 27 28 private static void stopServer(List<KafkaConsumerRunner> consumers) { 29 System.out.println("start shutdown"); 30 for(KafkaConsumerRunner consumerRunner : consumers){ 31 consumerRunner.shutdown(); 32 } 33 } 34 35 private static List<KafkaConsumerRunner> startServer(Properties props, int totleThread, CountDownLatch countDown) { 36 ExecutorService threadPool = Executors.newFixedThreadPool(totleThread); 37 List<KafkaConsumerRunner> consumers = new ArrayList<KafkaConsumerRunner>(); 38 39 for (int i = 0; i < 5; i++) { 40 KafkaConsumerRunner consumer = new KafkaConsumerRunner(countDown, props); 41 threadPool.execute(consumer); 42 consumers.add(consumer); 43 } 44 return consumers; 45 } 46 47 private final AtomicBoolean closed = new AtomicBoolean(false); 48 private final KafkaConsumer consumer = null; 49 private CountDownLatch countDown; 50 private Properties props; 51 private KafkaConsumer<String, String> consumer = null; 52 53 54 public KafkaConsumerRunner(CountDownLatch countDown, Properties props) { 55 this.countDown = countDown; 56 this.props = props; 57 this.consumer = new KafkaConsumer(this.props); 58 } 59 public void run() { 60 try { 61 consumer.subscribe(Arrays.asList("topic1", "demo")); 62 while (!closed.get()) { 63 ConsumerRecords<String, String> records = consumer.poll(Long.MAX_VALUE); 64 for (TopicPartition partition : records.partitions()) { 65 List<ConsumerRecord<String, String>> partitionRecords = records.records(partition); 66 for (ConsumerRecord<String, String> record : partitionRecords) { 67 // // Handle new records 68 System.out.println(record.offset() + ": --" + record.value() + " " + Thread.currentThread().getId()); 69 } 70 consumer.commitSync();//同步 71 } 72 } 73 } catch (WakeupException e) { 74 // Ignore exception if closing 75 if (!closed.get()) throw e; 76 } finally { 77 consumer.close(); 78 } 79 } 80 81 // Shutdown hook which can be called from a separate thread 82 public void shutdown() { 83 countDown.countDown(); 84 closed.set(true); 85 } 86 87}
一个多线程使用的例子
MockConsumer
这应该还是一个用来测试的实现
更多的配置参数可以看 http://kafka.apache.org/documentation#producerapi , 我们同样可以实现自己的 partition.class 只要实现 Partitioner 接口就可以
1public class RandomPartitioner implements Partitioner { 2 3 private Random ran = new Random(); 4 5 public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) { 6 synchronized (this) { 7 int partitionerNum = cluster.partitionsForTopic(topic).size(); 8 return ran.nextInt(partitionerNum); 9 } 10 } 11 12 public void close() { 13 14 } 15 16 public void configure(Map<String, ?> configs) { 17 18 } 19}
然后通过 Properties 类配置到 Producer 就可以了
props.put("partitioner.class", RandomPartitioner.class.getName());
3. Connectors
Connectors 是Kafka 0.9 版本提供的一种新的工具,他的主要目的是将其他系统(平台)的数据导入、导出到Kafka的集群里,甚至可以把整个数据库导入到Topic中,同时也可以作为数据的二次存储、查询或者批量的进行离线分析
Kafka 已经提供了一个工具给我们玩耍:
第一步:建立一个测试文件
echo "this is large data" > test.txt
第二步:查看source的配置文件
1vim config/connect-file-source.properties 2 3name=local-file-source 4connector.class=FileStreamSource 5tasks.max=1 6file=test.txt 7topic=connect-test
里面一个file参数是指定source(需要被导入kafka的数据)文件的位置
第三步:查看sink文件配置
1vim config/connect-file-sink.properties 2 3name=local-file-sink 4connector.class=FileStreamSink 5tasks.max=1 6file=test.sink.txt 7topics=connect-test
这里配置了数据从kafka集群导出后需要保存的位置(目前为test.sink.txt)
第四部:运行kafka自带脚本
bin/connect-standalone.sh config/connect-standalone.properties config/connect-file-source.properties config/connect-file-sink.properties
运行后会看到当前目录多了一个 text.sink.txt, 内容与 text.txt 一致,同时两个文件(数据)可以通过kafka保持同步
同时我们可以自己写一个配置文件链接数据库,例如
1name=test-mssql-jdbc-autoincrement 2connector.class=io.confluent.connect.jdbc.JdbcSourceConnector 3tasks.max=1 4connection.url=jdbc:sqlserver://127.0.0.1:1433;user=xd;password=xd;databaseName=scratchpad 5mode=incrementing 6incrementing.column.name=id 7topic.prefix=test-mssql-jdbc- 8table.whitelist=data01
同时我们可以通过实现 Connector 和 Task 接口来实现自己的应用,但很可惜我在客户端内没有找到这两个接口,只有在kafka源码中有这两个接口,官方的例子如下:
http://kafka.apache.org/documentation#connect\_developing
4. StreamProcessors
Stream是一个客户端库,用于处理和分析那些存储在Kafka上的数据,同时可以将结果写回到Kafka集群或者发送到外部的系统,其特点如下:
- 设计一个简单且轻量级的客户端(貌似Stream很多都是用spark),可以方便的嵌入到已有的Java应用当中
- 消息通讯层只依赖Kafka自己(不依赖其他东西),使用Kafka分区模型水平的拆分处理
- 支持本地状态容错(不是十分清楚什么鸟意思)
目前对Stream了解并不多,后续跟进