Kafka深入浅出(二)

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了解并不多,后续跟进

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )