Apache ECharts 介绍和使用

Apache ECharts
一个基于 JavaScript 的开源可视化图表库

5 分钟上手 ECharts

获取 ECharts

你可以通过以下几种方式获取 Apache EChartsTM。

引入 ECharts

通过标签方式直接引入构建好的 echarts 文件

<!DOCTYPE html>
<html>
<head>
    <meta charset="utf-8">
    <!-- 引入 ECharts 文件 -->
    <script type="text/javascript" src="https://cdn.jsdelivr.net/npm/echarts/dist/echarts.min.js"></script>
</head>
</html>

绘制一个简单的图表

在绘图前我们需要为 ECharts 准备一个具备高宽的 DOM 容器。

<body>
    <!-- 为 ECharts 准备一个具备大小(宽高)的 DOM -->
    <div id="main" style="width: 600px;height:400px;"></div>
</body>

然后就可以通过 echarts.init 方法初始化一个 echarts 实例并通过 setOption 方法生成一个简单的柱状图,下面是完整代码。

<!DOCTYPE html>
<html>
<head>
    <meta charset="utf-8">
    <title>ECharts</title>
    <!-- 引入 echarts.js -->
    <script src="echarts.min.js"></script>
</head>
<body>
    <!-- 为ECharts准备一个具备大小(宽高)的Dom -->
    <div id="main" style="width: 600px;height:400px;"></div>
    <script type="text/javascript">
        // 基于准备好的dom,初始化echarts实例
        var myChart = echarts.init(document.getElementById('main'));

        // 指定图表的配置项和数据
        var option = {
            title: {
                text: 'ECharts 入门示例'
            },
            tooltip: {},
            legend: {
                data:['销量']
            },
            xAxis: {
                data: ["衬衫","羊毛衫","雪纺衫","裤子","高跟鞋","袜子"]
            },
            yAxis: {},
            series: [{
                name: '销量',
                type: 'bar',
                data: [5, 20, 36, 10, 10, 20]
            }]
        };

        // 使用刚指定的配置项和数据显示图表。
        myChart.setOption(option);
    </script>
</body>
</html>

这样你的第一个图表就诞生了!

file

ECharts 基础概念概览

本文介绍 Apache EChartsTM 最基本的名词和概念。

echarts 实例

一个网页中可以创建多个 echarts 实例。每个 echarts 实例 中可以创建多个图表和坐标系等等(用 option 来描述)。准备一个 DOM 节点(作为 echarts 的渲染容器),就可以在上面创建一个 echarts 实例。每个 echarts 实例独占一个 DOM 节点。

img

系列(series)

系列series)是很常见的名词。在 echarts 里,系列series)是指:一组数值以及他们映射成的图。“系列”这个词原本可能来源于“一系列的数据”,而在 echarts 中取其扩展的概念,不仅表示数据,也表示数据映射成为的图。所以,一个 系列 包含的要素至少有:一组数值、图表类型(series.type)、以及其他的关于这些数据如何映射成图的参数。

echarts 里系列类型(series.type)就是图表类型。系列类型(series.type)至少有:line(折线图)、bar(柱状图)、pie(饼图)、scatter(散点图)、graph(关系图)、tree(树图)、...

如下图,右侧的 option 中声明了三个 系列series):pie(饼图系列)、line(折线图系列)、bar(柱状图系列),每个系列中有他所需要的数据(series.data)。

img

类同地,下图中是另一种配置方式,系列的数据从 dataset 中取:

img

组件(component)

在系列之上,echarts 中各种内容,被抽象为“组件”。例如,echarts 中至少有这些组件:xAxis(直角坐标系 X 轴)、yAxis(直角坐标系 Y 轴)、grid(直角坐标系底板)、angleAxis(极坐标系角度轴)、radiusAxis(极坐标系半径轴)、polar(极坐标系底板)、geo(地理坐标系)、dataZoom(数据区缩放组件)、visualMap(视觉映射组件)、tooltip(提示框组件)、toolbox(工具栏组件)、series(系列)、...

我们注意到,其实系列(series)也是一种组件,可以理解为:系列是专门绘制“图”的组件。

如下图,右侧的 option 中声明了各个组件(包括系列),各个组件就出现在图中。

img

注:因为系列是一种特殊的组件,所以有时候也会出现 “组件和系列” 这样的描述,这种语境下的 “组件” 是指:除了 “系列” 以外的其他组件。

用 option 描述图表

上面已经出现了 option 这个概念。echarts 的使用者,使用 option 来描述其对图表的各种需求,包括:有什么数据、要画什么图表、图表长什么样子、含有什么组件、组件能操作什么事情等等。简而言之,option 表述了:数据数据如何映射成图形交互行为

// 创建 echarts 实例。
var dom = document.getElementById('dom-id');
var chart = echarts.init(dom);

// 用 option 描述 <code>数据数据如何映射成图形交互行为 等。
// option 是个大的 JavaScript 对象。
var option = {
    // option 每个属性是一类组件。
    legend: {...},
    grid: {...},
    tooltip: {...},
    toolbox: {...},
    dataZoom: {...},
    visualMap: {...},
    // 如果有多个同类组件,那么就是个数组。例如这里有三个 X 轴。
    xAxis: [
        // 数组每项表示一个组件实例,用 type 描述“子类型”。
        {type: 'category', ...},
        {type: 'category', ...},
        {type: 'value', ...}
    ],
    yAxis: [{...}, {...}],
    // 这里有多个系列,也是构成一个数组。
    series: [
        // 每个系列,也有 type 描述“子类型”,即“图表类型”。
        {type: 'line', data: [['AA', 332], ['CC', 124], ['FF', 412], ... ]},
        {type: 'line', data: [2231, 1234, 552, ... ]},
        {type: 'line', data: [[4, 51], [8, 12], ... ]}
    }]
};

// 调用 setOption 将 option 输入 echarts,然后 echarts 渲染图表。
chart.setOption(option);

系列里的 series.data 是本系列的数据。而另一种描述方式,系列数据从 dataset 中取:

var option = {
    dataset: {
        source: [
            [121, 'XX', 442, 43.11],
            [663, 'ZZ', 311, 91.14],
            [913, 'ZZ', 312, 92.12],
            ...
        ]
    },
    xAxis: {},
    yAxis: {},
    series: [
        // 数据从 dataset 中取,encode 中的数值是 dataset.source 的维度 index (即第几列)
        {type: 'bar', encode: {x: 1, y: 0}},
        {type: 'bar', encode: {x: 1, y: 2}},
        {type: 'scatter', encode: {x: 1, y: 3}},
        ...
    ]
};

组件的定位

不同的组件、系列,常有不同的定位方式。

[类 CSS 的绝对定位]

多数组件和系列,都能够基于 top / right / down / left / width / height 绝对定位。 这种绝对定位的方式,类似于 CSS 的绝对定位(position: absolute)。绝对定位基于的是 echarts 容器 DOM 节点。

其中,他们每个值都可以是:

  • 绝对数值(例如 bottom: 54 表示:距离 echarts 容器底边界 54 像素)。
  • 或者基于 echarts 容器高宽的百分比(例如 right: '20%' 表示:距离 echarts 容器右边界的距离是 echarts 容器宽度的 20%)。

如下图的例子,对 grid 组件(也就是直角坐标系的底板)设置 leftrightheightbottom 达到的效果。

img

我们可以注意到,left right width 是一组(横向)、top bottom height 是另一组(纵向)。这两组没有什么关联。每组中,至多设置两项就可以了,第三项会被自动算出。例如,设置了 leftright 就可以了,width 会被自动算出。

[中心半径定位]

少数圆形的组件或系列,可以使用“中心半径定位”,例如,pie(饼图)、sunburst(旭日图)、polar(极坐标系)。

中心半径定位,往往依据 center(中心)、radius(半径)来决定位置。

[其他定位]

少数组件和系列可能有自己的特殊的定位方式。在他们的文档中会有说明。

坐标系

很多系列,例如 line(折线图)、bar(柱状图)、scatter(散点图)、heatmap(热力图)等等,需要运行在 “坐标系” 上。坐标系用于布局这些图,以及显示数据的刻度等等。例如 echarts 中至少支持这些坐标系:直角坐标系极坐标系地理坐标系(GEO)单轴坐标系日历坐标系 等。其他一些系列,例如 pie(饼图)、tree(树图)等等,并不依赖坐标系,能独立存在。还有一些图,例如 graph(关系图)等,既能独立存在,也能布局在坐标系中,依据用户的设定而来。

一个坐标系,可能由多个组件协作而成。我们以最常见的直角坐标系来举例。直角坐标系中,包括有 xAxis(直角坐标系 X 轴)、yAxis(直角坐标系 Y 轴)、grid(直角坐标系底板)三种组件。xAxisyAxisgrid 自动引用并组织起来,共同工作。

我们来看下图,这是最简单的使用直角坐标系的方式:只声明了 xAxisyAxis 和一个 scatter(散点图系列),echarts 暗自为他们创建了 grid 并关联起他们:

img

再来看下图,两个 yAxis,共享了一个 xAxis。两个 series,也共享了这个 xAxis,但是分别使用不同的 yAxis,使用 yAxisIndex 来指定它自己使用的是哪个 yAxis

img

再来看下图,一个 echarts 实例中,有多个 grid,每个 grid 分别有 xAxisyAxis,他们使用 xAxisIndexyAxisIndexgridIndex 来指定引用关系:

img

另外,一个系列,往往能运行在不同的坐标系中。例如,一个 scatter(散点图)能运行在 直角坐标系极坐标系地理坐标系(GEO) 等各种坐标系中。同样,一个坐标系,也能承载不同的系列,如上面出现的各种例子,直角坐标系 里承载了 line(折线图)、bar(柱状图)等等。

ECharts 高级

ECharts API文档

ECharts Option说明

ECharts 官方示例

Views: 529

Kakfa管理操作的相关命令

主题操作

kafka-topics.sh脚本可以用来管理主题.

--topic 指定操作的主题名,除了创建之外都可以使用正则表达式,但是要使用\转义。另外命名不要使用两个下划线开头(__表示系统内建的主题,比如__consumer_offset),主题命名中不要将.和_混用(kafka会将.最终转换成_)。

创建主题

kafka-topics.sh --bootstrap-server hadoop000:9092 --create --topic test-1 --partitions 2 --replication-factor 2

注意: --bootstrap-server 后不要使用localhost, 而是跟上主机名

kafka-topics.sh --bootstrap-server hadoop000:9092 --create --topic test-2 --partitions 2 --replication-factor 2

--zookeeper localhost:2181/kafka 这样的写法从Kafka 2.2版本开始废弃,新版本kafka推荐使用--bootstrap-server替换

其中localhost:2181/kafka 来自于 kafka的server.propertieszookeeper.connect的配置

描述主题

列出所有主题详情

kafka-topics.sh --bootstrap-server hadoop000:9092 --describe

支持正则表达式

 ~]$ kafka-topics.sh --bootstrap-server hadoop000:9092 --describe --topic test-\.\+

修改主题

alter不能修改副本因子。分区数只能增加,不能减少

kafka-topics.sh --bootstrap-server hadoop000:9092 --alter --topic test-1 --partitions 3

删除主题

--delete 只是标记主题为删除状态, 并不会马上删除, 如果你希望立即从zookeeper中删除此主题的信息, 可以在配置文件 server.properties 添加配置 delete.topic.enable=true

kafka-topics.sh --bootstrap-server hadoop000:9092 --delete --topic test-2

列出主题

只列出主题名称, 不带详细信息

kafka-topics.sh --bootstrap-server hadoop000:9092 --list

消费者操作

使用 kafka-consumer-groups.sh 工具,我们可以列出、描述或删除消费者组。消费者组可以被手动删除,也可以在该组最后一次提交的偏移量到期时被自动删除。手动删除仅在组中没有任何活动成员时有效。

列出消费者群组

> bin/kafka-consumer-groups.sh --bootstrap-server hadoop000:9092 --list
test-consumer-group-1
test-consumer-group-2
test-consumer-group-3

如果使用新的消费者客户端,要列出消费者组,可以使用--bootstrap-server和--list选项。因为新的客户端已经删除了--zookeeper选项。

描述消费者群组

> bin/kafka-consumer-groups.sh --bootstrap-server hadoop000:9092 
--describe --group my-group

TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENTID
topic1 0  854144 855809 1665 consumer1-3fc8d6f1-581a-4472-bdf3-3515b4aee8c1 /127.0.0.1 consumer1
topic2 0  460537  803290  342753  consumer1-3fc8d6f1-581a-4472-bdf3-3515b4aee8c1 /127.0.0.1 consumer1 
topic3 2  243655  398812  155157  consumer4-117fe4d3-c6c1-4178-8ee9-eb4a3954bee0 /127.0.0.1 consumer4

删除消费者群组

> bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 
--delete --group my-group --group my-other-group
Deletion of requested consumer groups('my-group', 'my-other-group') was successful.

当要删除的消费者组不为空时,执行上面命令你将得到以下错误: GroupNotEmptyException

消费偏移量管理

删除偏移量

删除某个主题的某个消费者组的偏移量

> kafka-consumer-groups.sh --bootstrap-server hadoop000:9092 --delete-offsets --group my-group --topic my-topic-1 --topic my-topic-2
TOPIC                          PARTITION       STATUS 
my-topic-1                     0               Successful
my-topic-2                     0               Successful

○当要删除的消费者组不为空时,执行上面命令你将得到以下错误: GroupNotEmptyException

重置偏移量

要重置消费者组的偏移量,可以使用“--reset-offsets”选项。此选项一次只支持一个消费者组。它需要定义以下范围: --all-topic 或 --topic。必须选择一个作用域,除非使用 “--from-file”从文件导入的方式。

> bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --reset-offsets --group consumergroup1 --topic topic1 --to-latest
TOPIC                          PARTITION  NEW-OFFSET
topic1                         0 

如果您使用的是旧的高级消费者,也就是消费者组元数据是存储在ZooKeeper中(即配置了offset .storage=ZooKeeper),则需要使用--zooKeeper而不是--bootstrap-server

配置项 描述
–-to-datetime 重置为某个时间点的偏移量.格式: 'YYYY-MM-DDTHH:mm:SS.sss'
--to-earliest : 重置为最早的偏移量
--to-latest : 重置为最近的偏移量
--shift-by : 将当前的偏移量偏移n个单位, n可以为正数也可以是负数
--from-file : 利用CSV文件中的数据重置偏移量
--to-current : R将偏移量重置为当前值
--by-duration : 重置偏移量为从当前时间戳开始的时长, 格式: 'PnDTnHnMnS'
--to-offset : 重置偏移量为指定的值

动态配置修改

从Kafka 1.1.0 开始有了动态配置修改的新特性。动态意味着在修改了Broker的配置后,我们不需要重新启动Broker来使其生效。

这个新特性包含在一个名为kafka-configs.sh的命令行工具脚本中。一旦使用这个工具设置了配置参数,新的更改将永久存储在zookeeper集群中。

覆盖主题默认配置

有许多应用于主题的配置,可以针对单个主题更改这些配置以适应集群中的不同用例。大多数配置都在代理配置中指定了缺省值,除非设置配置覆盖,否则将应用该缺省值。

格式

kafka-configs.sh --bootstrap-server <broker_list> --alter --entity-type topics --entity-name <topic name> --add-config <key>= <value>[,<key>=<value>...]

例子

# 下面的示例将my-topic主题的留存时间设置为1小时,即3600000ms:
$ kafka-configs.sh --bootstrap-server hadoop000:9092 --alter --entity-type topics --entity-name my-topic --add-config retention.ms=3600000

主题的有效可覆盖配置

配置(Key 说明
cleanup.policy 如果设置为compact,则topic中 的消息将被丢弃,仅保留具有给定key的最新消息(日志压缩)。
compression.type broker将消息写入磁盘时使用的压缩类型,可以用gzip、snappy和lz4.
delete.retention.ms 压缩日志墓碑消息的最长存放时间
file.delete.delay.ms 从磁盘中删除此topic的日志端和索引之前需要等待的多长时间
配置(Key 说明
flush.messages 在强制将此topic的消息刷到磁盘之前接收的消息数
flush.ms 在强制将此topic的消息刷到磁盘之前需要的时间,单位是毫秒
index.interval.bytes 日志段索引中的条目之间可以产生多少字节的消息
max.message.bytes 此topic中当个消息的大小
retention.bytes 为topic保留的消息量的总字节数
retention.ms topic中消息保留的最长时间,单位是毫秒

描述覆盖的配置

可以使用命令行工具kafka-configs.sh来检查主题或客户机的特定配置。显示覆盖的配置需要使用--describe选项。例如,一下代码显示名为my-topic的主题的所有覆盖过的配置

> kafka-configs.sh --bootstrap-server <broker_list> --describe --entity-type topics --entity-name my-topic

Configs for topic 'my-topic' are retention.ms=3600000

覆盖客户端默认配置

Kafka客户端(生产者和消费者)唯一可以覆盖的配置是生产者和消费者的配额,即允许具有指定客户端ID的所有客户端在每个broker上每秒生产或者消费的字节数。

格式

> kafka-configs.sh --bootstrap-server <broker_list> --alter --entity-type clients --entity-name <client ID> --add-config <key>=<value>[,<key>=<value>...]

Kafka客户端的有效配置

配置key 描述
producer_bytes_rate The amount of messages, in bytes, that a singe client ID is allowed to produce to a single broker in one second.
consumer_bytes_rate The amount of messages, in bytes, that a single client ID is allowed to consume from a single broker in one second.

如何把Broker中的生产者客户端ID为client_a的配额producer_byte_rate设置成20MB/秒?

> bin/kafka-configs.sh
 --bootstrap-server hadoop000:9092 
 --alter
 --entity-type clients 
 --entity-name client_a
 --add-config 'producer_bytes_rate=20971520'

注意:

  • --add-config 可以同时指定多个配置, 配置之间使用逗号分割

  • 要指定生产者或者消费者的client_id,如果是Java APIs可以为客户端添加配置client.id, 如果是控制台生产者或者消费者可以使用--producer.config--consumer.config进行指定配置文件并且在配置文件中指定client.id.

移除覆盖的配置

可以完全删除动态配置,这将导致集群配置恢复到默认值,要删除配置覆盖,请使用--alter命令以及--delete-config命令。

下面的示例可以删除一个名为my-topic的主题的覆盖后的retention.ms配置, 删除后retention.ms将恢复为默认值:

kafka-configs.sh --bootstrap-server hadoop000:9092 --alter --entity-type topics --entity-name my-topic --delete-config retention.ms

Completed updating config for entity: topic 'my-topic'.

Views: 641

Storm – 使用Trident实现词频统计并提供实时查询

为什么使用Trident

逐个处理单个tuple会增加很多开销,因此storm中引入Trident实现batch处理.

Trident优点是:

  • 批次处理消息
  • 减少持久化的开销
  • 结合Trident State能可靠保证每个消息只被处理一次

Trident的 State

Trident 在进行聚合操作时需要缓存中间结果, 可以看做Trident的状态(State).
Trident状态既可以保留在topology的内部,比如说内存中,也可以放到外部存储当中,比如说Memcached或者Cassandra数据库中.

Trident 允许以一种容错的方式来管理状态, 可以保证当遇到重试或错误时状态的更新是幂等的, 以此来实现EOS(Exactly only semantics)。

注:在数据统计分析中,幂等性是一个很重要的指标,因为它可以保证即使数据被处理了多次,但是站在结果的角度看和处理一次完全一样。

使用Trident实现词频统计

这里我们词频统计使用Trident拓扑实现.
这里使用Storm提供的MemoryMapState管理状态(每个单词的实时统计个数), 顾名思义MemoryMapState管理的状态是保存在内存中的. 状态可以通过方法stateQuery进行实时查询.

然后提供一个DRPC服务, 客户端可以指定需要查询词频的单词有哪些.

本地提交模式

数据源

在这个例子中,我们使用FixedBatchSpout 对象模拟数据源来发送一句一句的文本内容。输入数据源也可以和Kestrel或者Kafka这样的消息队列对接, Trident在处理输入流的时候会转换成若干个tuple组成的batch来处理。这个 FixedBatchSpoutmaxBatchSize为5.

FixedBatchSpout spout = new FixedBatchSpout(new Fields("sentence"), 5,
                new Values("the cow jumped over the moon"),
                new Values("the man went to the store and bought some candy"),
                new Values("four score and seven years ago"),
                new Values("how many apples can you eat"),
                new Values("to be or not to be the person"));
        spout.setCycle(false);

Trident实现词频统计

接下来我们创建TridentTopology计算词频, 并将实时统计的结果保存在TridentState中:

TridentTopology topology = new TridentTopology();

        Config conf = new Config();
        conf.setMaxSpoutPending(20);
        conf.setNumWorkers(3);

        LocalDRPC drpc = new LocalDRPC();
        LocalCluster localCluster = new LocalCluster();

        TridentState wordCounts =
                topology.newStream("spout1", spout).parallelismHint(16)
                        .each(new Fields("sentence"), new Split(), new Fields("word"))
                        .groupBy(new Fields("word"))
                        .persistentAggregate(new MemoryMapState.Factory(), new Count(), new Fields("count"))
                        .parallelismHint(3);

解析:

  • TridentTopologynewStream方法传入了一个spout对象,spout对象会从外部读取数据并输出到当前topology当中,从而在topology中创建了一个新的数据流.
  • Trident会在Zookeeper中保存一小部分状态信息来追踪数据的处理情况,而在代码中我们指定的字符串“spout1”就是Zookeeper中用来存储metadata信息的Znode节点.
  • persistentAggregate方法会把数据流转换成一个TridentState对象, TridentState记录了单词的实时词频。

处理过程:

  • 根据空格拆分sentence,并将拆分出的每个单词作为一个tuple输出
  • 根据“word”字段进行groupBy操作
  • persistentAggregate会帮助你把计数的结果进行存储

提供DRPC服务实现实时词频查询

通过newDRPCStream可以接受DRPC客户端请求的参数:

topology.newDRPCStream("words", drpc)
                .each(new Fields("args"), new Split(), new Fields("word"))
                .groupBy(new Fields("word"))
                .stateQuery(wordCounts, new Fields("word"), new MapGet(), new Fields("count"))
                .each(new Fields("count"), new FilterNull())
                .project(new Fields("word", "count"))
                .aggregate(new Fields("count"), new Sum(), new Fields("sum"));

解析:

  • 对参数按照空格切分后作为word字段发射出去, 再对word 进行groupBy操作.
  • 使用stateQuery来在上面代码中创建的TridentState对象上进行查询。
  • stateQuery利用MapGet来获取每个单词的出现个数。
  • 由于DRPC stream是使用跟TridentState完全同样的group方式(按照“word”字段进行group),每个单词的查询会被路由到TridentState的相应分区去执行。
  • FilterNull这个过滤器把从未出现过的单词给去掉,
  • 并使用Sum这个聚合器将这些词频统计结果累加起来。最终,Trident会自动把这个结果发送回等待的客户端。

提交到集群

为了更快看到结果, 这里使用本地提交方式:

        localCluster.submitTopology(topoName, conf, topology.build());

DRPC本地客户端发起实时查询

下面这部分实现了一个低延时的单词数量的分布式实时DRPC查询。这个查询以一个用空格分割的单词列表为输入,并返回这些单词事实出现次数。

这些查询是像普通的RPC调用那样被执行的,要说不同的话,那就是他们在后台是并行执行的。下面是执行DRPC实时查询的一个例子:

        for (int i = 0; i < 10; i++) {
            System.out.println("DRPC RESULT: " + drpc.execute("words", "cat the dog jumped"));
            Thread.sleep(1000);
        }

        localCluster.shutdown();
        drpc.shutdown();

解析:

  • 发起DRPC请求, 请求的参数是"cat the dog jumped", 每隔1秒查询一次, 共10次.

Trident在如何最大程度的保证执行topogloy性能方面是非常智能的。在topology中会自动的发生两件非常有意思的事情:

  • 更新状态和读取状态操作 (比如说 persistentAggregate 和 stateQuery) 会自动的是batch的形式操作状态。
  • 如果有20次更新需要被同步到存储中,Trident会自动的把这些操作汇总到一起批处理,只做一次读写操作,而不是进行20次读写操作。

因此你可以在很方便的执行计算的同时,保证了非常好的性能。

Trident 的聚合器已经是被优化的非常好了的。Trident并不是简单的把一个group中所有的tuples都发送到同一个机器上面进行聚合,而是在发送之前已经进行过一次局部的聚合。打个比方,Count聚合器会先在每个partition上面进行count,然后把每个分片count汇总到一起就得到了最终 的count。这个技术其实就跟MapReduce里面的combiner是一个思想。

完整代码代码:

import org.apache.storm.Config;
import org.apache.storm.LocalCluster;
import org.apache.storm.LocalDRPC;
import org.apache.storm.trident.TridentState;
import org.apache.storm.trident.TridentTopology;
import org.apache.storm.trident.operation.BaseFunction;
import org.apache.storm.trident.operation.TridentCollector;
import org.apache.storm.trident.operation.builtin.Count;
import org.apache.storm.trident.operation.builtin.FilterNull;
import org.apache.storm.trident.operation.builtin.MapGet;
import org.apache.storm.trident.testing.FixedBatchSpout;
import org.apache.storm.trident.testing.MemoryMapState;
import org.apache.storm.trident.tuple.TridentTuple;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Values;

/**
 * 将wordcount结果持久化到数据库
 */
public class AxTridentWordCountLocally {

    public static void main(String[] args) throws Exception {
        String topoName = "wordCounter";

        FixedBatchSpout spout = new FixedBatchSpout(new Fields("sentence"), 5,
                new Values("the cow jumped over the moon"),
                new Values("the man went to the store and bought some candy"),
                new Values("four score and seven years ago"),
                new Values("how many apples can you eat"),
                new Values("to be or not to be the person"));
        spout.setCycle(false);

        TridentTopology topology = new TridentTopology();

        Config conf = new Config();
        conf.setMaxSpoutPending(20);
        conf.setNumWorkers(3);

        LocalDRPC localDRPC = new LocalDRPC();
        LocalCluster localCluster = new LocalCluster();

        /*
         * 根据数据源拆分单词后,然后分区操作,在每个分区上又进行分组(hash算法),然后在每个分组上进行聚合
         * 所以这里可能有多个分区,每个分区有多个分组,然后在多个分组上进行聚合
         * 用来进行group的字段会以key的形式存在于State当中,聚合后的结果会以value的形式存储在State当中
         */
        TridentState wordCounts = topology.newStream("spout1", spout).parallelismHint(16)
                .each(new Fields("sentence"), new Split(), new Fields("word"))
                .groupBy(new Fields("word"))
                .persistentAggregate(new MemoryMapState.Factory(), new Count(), new Fields("count"))
                .parallelismHint(3);

        topology.newDRPCStream("words", localDRPC).each(new Fields("args"), new Split(), new Fields("word"))
                .groupBy(new Fields("word"))
                .stateQuery(wordCounts, new Fields("word"), new MapGet(), new Fields("count"))
                .each(new Fields("count"), new FilterNull())
                .project(new Fields("word", "count"));

        localCluster.submitTopology(topoName, conf, topology.build()); //异步

        for (int i = 0; i < 10; i++) {
            System.out.println("DRPC RESULT: " + localDRPC.execute("words", "cat the dog jumped"));
            Thread.sleep(1000);
        }

        localCluster.shutdown();
        localDRPC.shutdown();
    }

    public static class Split extends BaseFunction {
        @Override
        public void execute(TridentTuple tuple, TridentCollector collector) {
            String sentence = tuple.getString(0);
            for (String word : sentence.split(" ")) {
                collector.emit(new Values(word));
            }
        }
    }
}

执行结果

在执行上述代码之后,可能输出如下所示:

DRPC RESULT: []
DRPC RESULT: [["the",30],["jumped",6]]
DRPC RESULT: [["the",84],["jumped",16]]
DRPC RESULT: [["jumped",26],["the",130]]
DRPC RESULT: [["jumped",30],["the",149]]
DRPC RESULT: [["jumped",40],["the",199]]
DRPC RESULT: [["the",245],["jumped",49]]
DRPC RESULT: [["jumped",59],["the",295]]
DRPC RESULT: [["the",345],["jumped",69]]
DRPC RESULT: [["the",394],["jumped",79]]

数据流向示意图

file

远程提交模式

import org.apache.storm.Config;
import org.apache.storm.StormSubmitter;
import org.apache.storm.generated.StormTopology;
import org.apache.storm.trident.TridentState;
import org.apache.storm.trident.TridentTopology;
import org.apache.storm.trident.operation.BaseFunction;
import org.apache.storm.trident.operation.TridentCollector;
import org.apache.storm.trident.operation.builtin.Count;
import org.apache.storm.trident.operation.builtin.FilterNull;
import org.apache.storm.trident.operation.builtin.MapGet;
import org.apache.storm.trident.testing.FixedBatchSpout;
import org.apache.storm.trident.testing.MemoryMapState;
import org.apache.storm.trident.tuple.TridentTuple;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Values;
import org.apache.storm.utils.DRPCClient;

public class CxTridentWordCountRemotely {
    public static StormTopology buildTopology() {
        FixedBatchSpout spout = new FixedBatchSpout(new Fields("sentence"), 3, new Values("the cow jumped over the moon"),
                new Values("the man went to the store and bought some candy"),
                new Values("four score and seven years ago"),
                new Values("how many apples can you eat"), new Values("to be or not to be the person"));
        spout.setCycle(true);

        TridentTopology topology = new TridentTopology();

        TridentState wordCounts = topology.newStream("spout1", spout).parallelismHint(16).each(new Fields("sentence"),
                new Split(), new Fields("word"))
                .groupBy(new Fields("word")).persistentAggregate(new MemoryMapState.Factory(),
                        new Count(), new Fields("count"))
                .parallelismHint(16);

        topology.newDRPCStream("words").each(new Fields("args"), new Split(), new Fields("word"))
                .groupBy(new Fields("word"))
                .stateQuery(wordCounts, new Fields("word"), new MapGet(), new Fields("count"))
                .each(new Fields("count"), new FilterNull())
                .project(new Fields("word", "count"));
        return topology.build();
    }

    public static void main(String[] args) throws Exception {
        Config conf = new Config();
        conf.setMaxSpoutPending(20);
        String topoName = "wordCounter";
        if (args.length > 0) {
            topoName = args[0];
        }
        conf.setNumWorkers(3);
        StormSubmitter.submitTopologyWithProgressBar(topoName, conf, buildTopology());
        try (DRPCClient drpc = DRPCClient.getConfiguredClient(conf)) {
            for (int i = 0; i < 10; i++) {
                System.out.println("DRPC RESULT: " + drpc.execute("words", "cat the dog jumped"));
                Thread.sleep(1000);
            }
        }
    }
    public static class Split extends BaseFunction {
        @Override
        public void execute(TridentTuple tuple, TridentCollector collector) {
            String sentence = tuple.getString(0);
            for (String word : sentence.split(" ")) {
                collector.emit(new Values(word));
            }
        }
    }
}

打成jar包上传到集群

Views: 628

Kafka中的消息序列化和反序列化

Kafka生产者中的配置项key.serializervalue.serializer指示如何将用户通过其ProducerRecord提供的键和值对象转换为字节。对于简单的字符串或字节类型,可以使用包含的ByteArraySerializerStringSerializer进行序列化操作。

kafka在发送或者接收消息的时候实际是使用byte[]字节型数组进行传输的。但是我们平常使用的时候,不但可以使用byte[],还可以使用int、short、long、float、double、String等数据类型,这是因为在我们使用这些数据类型的时候,kafka根据我们指定的序列化和反序列化方式转成byte[]类型之后再进行传输来提高传输效率。

通常我们在使用kakfa发送或者接受消息的时候都需要指定消息的key和value序列化方式,如生产者我们可以设置value.serializerorg.apache.kafka.common.serialization.StringSerializer来设置value的序列化方式为字符串,即我们可以发送string类型的消息。目前kafka原生支持的序列化和反序列化方式如下两表所示,这些原生的序列化和反序列化的类都是在org.apache.kafka.common.serialization包之下:

kafka序列化方式表

序列化方式 对应java数据类型 说明
ByteArraySerializer byte[] 原生类型
ByteBufferSerializer ByteBuffer 关于ByteBuffer
IntegerSerializer Interger
ShortSerializer Short
LongSerializer Long
DoubleSerializer Double
StringSerializer String

kafka反序列化方式表

序列化方式 对应java数据类型 说明
ByteArrayDeserializer byte[] 原生类型
ByteBufferDeserializer ByteBuffer 关于ByteBuffer
IntegerDeserializer Interger
ShortDeserializer Short
LongDeserializer Long
DoubleDeserializer Double
StringDeserializer String

Java原生类型的序列化和反序列化

上面我们了解一些关于kafka原生的一些序列化和反序列化方式。它们究竟是如实现的呢?以string类型为例子,我们看一下,kafka如何实现序列化/反序列化的。

kafka序列化/反序列化方式的实现代码在org.apache.kafka.common.serialization包下。

String 类型的序列化类

我们查看org.apache.kafka.common.serialization.StringSerializer这个类。

package org.apache.kafka.common.serialization;

import org.apache.kafka.common.errors.SerializationException;

import java.io.UnsupportedEncodingException;
import java.util.Map;

/**
 *  String encoding defaults to UTF8 and can be customized by setting the property key.serializer.encoding,
 *  value.serializer.encoding or serializer.encoding. The first two take precedence over the last.
 */
public class StringSerializer implements Serializer<String> {
    private String encoding = "UTF8";

    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
        String propertyName = isKey ? "key.serializer.encoding" : "value.serializer.encoding";
        Object encodingValue = configs.get(propertyName);
        if (encodingValue == null)
            encodingValue = configs.get("serializer.encoding");
        if (encodingValue instanceof String)
            encoding = (String) encodingValue;
    }

    @Override
    public byte[] serialize(String topic, String data) {
        try {
            if (data == null)
                return null;
            else
                return data.getBytes(encoding);
        } catch (UnsupportedEncodingException e) {
            throw new SerializationException("Error when serializing string to byte[] due to unsupported encoding " + encoding);
        }
    }

    @Override
    public void close() {
        // nothing to do
    }
}

由上面的代码我们可以看出:

  • String的序列化类是继承了Serializer接口,指定<String>泛型,然后实现的Serializer接口的configure()serialize()close()方法。代码重点的实现是在serialize(),可以看出这个方法将我们传入的String类型的数据,简单的通过data.getBytes(encoding)方法进行了序列化。

String 类新的反序列化类

我们查看org.apache.kafka.common.serialization.StringDeserializer这个类。

package org.apache.kafka.common.serialization;

import org.apache.kafka.common.errors.SerializationException;

import java.io.UnsupportedEncodingException;
import java.util.Map;

/**
 *  String encoding defaults to UTF8 and can be customized by setting the property key.deserializer.encoding,
 *  value.deserializer.encoding or deserializer.encoding. The first two take precedence over the last.
 */
public class StringDeserializer implements Deserializer<String> {
    private String encoding = "UTF8";

    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
        String propertyName = isKey ? "key.deserializer.encoding" : "value.deserializer.encoding";
        Object encodingValue = configs.get(propertyName);
        if (encodingValue == null)
            encodingValue = configs.get("deserializer.encoding");
        if (encodingValue instanceof String)
            encoding = (String) encodingValue;
    }

    @Override
    public String deserialize(String topic, byte[] data) {
        try {
            if (data == null)
                return null;
            else
                return new String(data, encoding);
        } catch (UnsupportedEncodingException e) {
            throw new SerializationException("Error when deserializing byte[] to string due to unsupported encoding " + encoding);
        }
    }

    @Override
    public void close() {
        // nothing to do
    }
}

同样,由上面的代码我们可以看出:

  • String类型的反序列化类是继承了Deserializer接口,指定<String>泛型,然后实现的Deserializer接口的configure()deserialize()close()方法。代码重点的实现是在deserialize(),可以看出这个方法将我们传入的byte[]类型的数据,简单的通过return new String(data, encoding)方法进行了反序列化得到了String类型的数据。

复杂Java对象的序列化及反序列化

通过上面,我们对kafka原生序列化/反序列化方式的了解,我们可以看出,kafka实现序列化/反序列化可以简单的总结为两步,第一步实现序列化Serializer或者反序列化Deserializer接口。第二步实现接口方法,将指定类型序列化成byte[]或者将byte[]反序列化成指定数据类型(String)。

graph LR
    生产者(String)--序列化-->传输和持久("Byte[]")--反序列化-->消费者(String)

由于原生类型不支持复杂对象类型,所以接下来,我们来实现对复杂对象类型的自定义序列化/反序列化方式。

这里我们介绍两种方式:

利用Buffer缓冲区

问题陈述:

Chapter 4 Activity 4.2

Kafka 消息也可以传递复杂对象类型。通常,这些对象将有多个字段。例如供应商对象。

实体类 Supplier

@Data
public class Supplier {
    private final int supplierId;
    private final String supplierName;
}

如果希望发送此类自定义对象或类型结构,则需要实现自定义序列化器和反序列化器。实现需要自定义序列化器和反序列化器。

解决方案 为了实现上述要求,利用缓冲区(Buffer,内存中预留指定大小的存储空间)用来对输入/输出(I/O)的数据作临时存储,这里由于需要保存字节数据所以我们使用ByteBuffer

自定义序列化类 SupplierSerializer

public class SupplierSerializer implements Serializer<Supplier> {
    private final String encoding = "UTF8";

    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {

    }
    @Override
    public byte[] serialize(String topic, Supplier data) {

        int sizeOfName;
        byte[] serializedName;

        try {
            if (data == null) {
                return null;
            }
           // 根据 字符编码 转换成字节数组
            serializedName = data.getName().getBytes(encoding);
            // 通过sizeOfName可以帮助在反序列化的时候
            // 可以知道要取多少个字符来根据指定编码格式转换成回字符串
            sizeOfName = serializedName.length;
            ByteBuffer buf = ByteBuffer.allocate(4 + 4 + sizeOfName);

            buf.putInt(data.getID()); // supplierId 是int类型占4 bytes
            buf.putInt(sizeOfName);   // supplierName 长度
            buf.put(serializedName);  // supplierName 内容

            return buf.array();

        } catch (Exception e) {
            throw new SerializationException("Error when serializing Supplier to byte[]");
        }
    }

    @Override
    public void close() {
        // nothing to do
    }
}

自定义反序列化类 SupplierDeserializer

public class SupplierDeserializer implements Deserializer<Supplier> {
    private final String encoding = "UTF8";

    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
        //Nothing to configure
    }

    @Override
    public Supplier deserialize(String topic, byte[] data) {

        try {
            if (data == null) {
                System.out.println("Null recieved at deserialize");
                return null;
            }

            ByteBuffer buf = ByteBuffer.wrap(data);
            int id = buf.getInt(); // 获取 supplierId的 int类型的值

            // 获取 supplierName长度
            int sizeOfName = buf.getInt();
            // 创建相同长度的字节数组
            byte[] nameBytes = new byte[sizeOfName]; 

            // get(nameBytes) :
            //  从当前位置开始相对读,读nameBytes.length个byte,
            //  并写入dst下标从offset到offset+length的区域
            buf.get(nameBytes); // 根据获取 supplierName
            String deserializedName = new String(nameBytes, encoding);

            return new Supplier(id, deserializedName);

        } catch (Exception e) {
            throw new SerializationException("Error when deserializing byte[] to Supplier");
        }
    }

    @Override
    public void close() {
        // nothing to do
    }
}

JSON格式的序列化及反序列化

实体类 JsonData

import java.util.Date;
import java.util.Map;

// lombok
@Data
@AllArgsConstructor
@NoArgsConstructor
public class JsonData {
    @JsonProperty("lng")
    private double longitude;
    @JsonProperty("lat")
    private double latitude;
    private double weight;
    private Date timestamp;
}

什么是 JSON ?

  • JSON 指的是 JavaScript 对象表示法(JavaScript Object Notation)
  • JSON 是轻量级的文本数据交换格式,类似 XML
  • JSON 比 XML 更小、更快,更易解析。
  • JSON 独立于语言:JSON 使用 Javascript语法来描述数据对象,但是 JSON 仍然独立于语言和平台。JSON 解析器和 JSON 库支持许多不同的编程语言。 目前非常多的动态(PHP,JSP,.NET)编程语言都支持JSON。
  • JSON 具有自我描述性,更易理解

JSON 实例

{
    "sites": [    
        { "name":"NIIT" , "url":"www.niit.com.cn" },     
        { "name":"Google" , "url":"www.google.com" },     
        { "name":"WeiBo" , "url":"www.weibo.com" }    
    ]
}

这个 sites 对象是包含 3 个站点记录(对象)的数组。

语法:

JSON 键必须是字符串,字符串必须使用双引号包裹。

JSON 值可以是:

  • 数字(整数或浮点数)
  • 字符串(在双引号中)
  • 逻辑值(true 或 false)
  • 数组(在中括号中)
  • 对象(在大括号中)
  • null

JSON的Java解析库 - Jackson

市面有很多用于解析JSON的用Java编写的第三方库,比如Jackson(fasterxml,可靠、灵活、可定制,使用广泛),Gson(Google, 轻量、简洁),国内比较著名的有FastJson(Alibaba个人开源,特点是快,虽然很有多历史漏洞,但是作为国人还是要支持一下,可以越来越好),这些库的用法比较类似,这里的案例使用的是Jackson

Jackson的maven依赖

<dependency>
    <groupId>com.fasterxml.jackson.core</groupId>
    <artifactId>jackson-core</artifactId>
    <version>2.11.0</version>
</dependency>
<dependency>
    <groupId>com.fasterxml.jackson.core</groupId>
    <artifactId>jackson-databind</artifactId>
    <version>2.11.0</version>
</dependency>
<dependency>
    <groupId>com.fasterxml.jackson.core</groupId>
    <artifactId>jackson-annotations</artifactId>
    <version>2.11.0</version>
</dependency>

基于Jackson的一个用于解析Json的工具类 JsonSerializerUtil

public class JsonSerializerUtil {

    /**
     * JSON序列化
     *
     * @param object 对象
     * @return JSON字符串
     */
    public static String serialize(Object object) {
        ObjectMapper mapper = new ObjectMapper();
        try {
            return mapper.writeValueAsString(object);
        } catch (JsonProcessingException e) {
            e.printStackTrace();
            return "";
        }
    }

    /**
     * JSON字符串反序列化
     *
     * @param jsonStr JSON字符串
     * @return a Map
     */
    public static Map deserialize(String jsonStr) {
        try {
            return deserialize(jsonStr, Map.class);
        } catch (Exception e) {
            e.printStackTrace();
            return new HashMap();
        }
    }

    public static <T> T deserialize(String jsonStr, Class<T> classType) throws Exception {
        return new ObjectMapper().readValue(jsonStr, classType);
    }
}

自定义序列化类 JsonDataSerializer

public class JsonDataSerializer implements Serializer<JsonData> {

    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {

    }

    @Override
    public byte[] serialize(String topic, JsonData data) {
        return JsonSerializerUtil.serialize(data).getBytes();
    }

    @Override
    public void close() {
        // nothing to do
    }
}

自定义反序列化类 JsonDataDeserializer

public class JsonDataDeserializer implements Deserializer<JsonData> {
    private final String encoding = "UTF8";

    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
        //Nothing to configure
    }

    @Override
    public JsonData deserialize(String topic, byte[] data) {

        if (data == null) {
            return null;
        }

        JsonData jsonData = null;

        try {
            String jsonString = new String(data, encoding);
            jsonData =  JsonSerializerUtil.deserialize(jsonString, JsonData.class);
        } catch (Exception e) {
            e.printStackTrace();
        }

        return jsonData;
    }

    @Override
    public void close() {
        // nothing to do
    }
}

作业

LG Activity Exer 1

需要将书籍(Book) 的信息通过Kafka传输, 书籍的相关字段有 书名 (name), 订购数量(quantityOrdered),单价(unitPrice),请编写一个生产者发送书籍消息, 在编写一个消费者消费书籍消息。

可以自行选择序列化的方式。

总结

实现序列化还有很多比较成熟的第三方序列化库可以使用(如avro,protoBuff等),关于采用什么样的方式去序列化数据还需要根据业务场景自己去定义。

Views: 569

Kafka Consumer API

高级API

在控制台创建发送者

$ bin/kafka-console-producer.sh \
--broker-list hadoop000:9092 --topic first

>hello world

创建消费者(过时API)

import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Properties;
import kafka.consumer.Consumer;
import kafka.consumer.ConsumerConfig;
import kafka.consumer.ConsumerIterator;
import kafka.consumer.KafkaStream;
import kafka.javaapi.consumer.ConsumerConnector;

public class CustomConsumer {

  @SuppressWarnings("deprecation")
  public static void main(String[] args) {
     Properties properties = new Properties();
     properties.put("zookeeper.connect", "hadoop000:2181");
     properties.put("group.id", "g1");
     properties.put("zookeeper.session.timeout.ms", "500");
     properties.put("zookeeper.sync.time.ms", "250");
     properties.put("auto.commit.interval.ms", "1000");

     // 创建消费者连接器
     ConsumerConnector consumer = Consumer.createJavaConsumerConnector(new ConsumerConfig(properties));

     HashMap<String, Integer> topicCount = new HashMap<>();
     topicCount.put("first", 1);

     Map<String, List<KafkaStream<byte[], byte[]>>> consumerMap = consumer.createMessageStreams(topicCount);

     KafkaStream<byte[], byte[]> stream = consumerMap.get("first").get(0);

     ConsumerIterator<byte[], byte[]> it = stream.iterator();

     while (it.hasNext()) {
       System.out.println(new String(it.next().message()));
     }
  }
}

官方提供案例(自动维护消费情况, 新API)

高级消费者和简单的消费者有以下的区别。

1.自动/隐藏偏移管理(Offset Management )

2.自动(简单)分区分配

3.Broker 故障转移 => 自动重新平衡

4.Consumer 故障转移 => 自动重新平衡

import java.util.Arrays;
import java.util.Properties;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;

public class CustomNewConsumer {

  public static void main(String[] args) {

     Properties props = new Properties();
     // 定义kakfa 服务的地址,不需要将所有broker指定上 
     props.put("bootstrap.servers", "hadoop000:9092");
     // 制定consumer group 
     props.put("group.id", "test");
     // 是否自动确认offset 
     props.put("enable.auto.commit", "true");
     // 自动确认offset的时间间隔 
     props.put("auto.commit.interval.ms", "1000");
     // key的序列化类
     props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
     // value的序列化类 
     props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
     // 定义consumer 
     KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
     // 消费者订阅的topic, 可同时订阅多个 
     consumer.subscribe(Arrays.asList("first", "second","third"));

     while (true) {
       // 读取数据,读取超时时间为100ms 
       ConsumerRecords<String, String> records = consumer.poll(100);

       for (ConsumerRecord<String, String> record : records)
         System.out.printf("offset = %d, key = %s, value = %s%n", record.offset(), record.key(), record.value());
     }
  }
}

低级API

也叫Simple Consumer, 实际使用起来并不简单.

实现使用低级API读取指定topic,指定partition,指定offset的数据。

1)消费者使用低级API 的主要步骤:

步骤 主要工作
1 根据指定的分区从主题元数据中找到主副本
2 获取分区最新的消费进度
3 从主副本拉取分区的消息
4 识别主副本的变化,重试

2)方法描述:

findLeader() 客户端向种子节点发送主题元数据,将副本集加入备用节点
getLastOffset() 消费者客户端发送偏移量请求,获取分区最近的偏移量
run() 消费者低级AP I拉取消息的主要方法
findNewLeader() 当分区的主副本节点发生故障,客户将要找出新的主副本

3)代码:

import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;

import kafka.api.FetchRequest;
import kafka.api.FetchRequestBuilder;
import kafka.api.PartitionOffsetRequestInfo;
import kafka.cluster.BrokerEndPoint;
import kafka.common.ErrorMapping;
import kafka.common.TopicAndPartition;
import kafka.javaapi.FetchResponse;
import kafka.javaapi.OffsetResponse;
import kafka.javaapi.PartitionMetadata;
import kafka.javaapi.TopicMetadata;
import kafka.javaapi.TopicMetadataRequest;
import kafka.javaapi.consumer.SimpleConsumer;
import kafka.message.MessageAndOffset;

public class SimpleExample {
  private List<String> m_replicaBrokers = new ArrayList<>();

  public SimpleExample() {
    m_replicaBrokers = new ArrayList<>();
  }

  public static void main(String args[]) {
    SimpleExample example = new SimpleExample();
    // 最大读取消息数量
    long maxReads = Long.parseLong("3");
    // 要订阅的topic
    String topic = "test1";
    // 要查找的分区
    int partition = Integer.parseInt("0");
    // broker节点的ip
    List<String> seeds = new ArrayList<>();
    seeds.add("192.168.9.102");
    seeds.add("192.168.9.103");
    seeds.add("192.168.9.104");
    // 端口
    int port = Integer.parseInt("9092");
    try {
      example.run(maxReads, topic, partition, seeds, port);
    } catch (Exception e) {
      System.out.println("Oops:" + e);
      e.printStackTrace();
    }
  }

  public void run(long a_maxReads, String a_topic, int a_partition, List<String> a_seedBrokers, int a_port) throws Exception {

    // 获取指定Topic partition的元数据
    PartitionMetadata metadata = findLeader(a_seedBrokers, a_port, a_topic, a_partition);

    if (metadata == null) {
      System.out.println("Can't find metadata for Topic and Partition. Exiting");
      return;
    }

    if (metadata.leader() == null) {
      System.out.println("Can't find Leader for Topic and Partition. Exiting");
      return;
    }

    String leadBroker = metadata.leader().host();
    String clientName = "Client" + a_topic + "" + a_partition;

    SimpleConsumer consumer = new SimpleConsumer(leadBroker, a_port, 100000, 64 * 1024, clientName);
    long readOffset = getLastOffset(consumer, a_topic, a_partition, kafka.api.OffsetRequest.EarliestTime(), clientName);
    int numErrors = 0;

    while (a_maxReads > 0) {
      if (consumer == null) {
        consumer = new SimpleConsumer(leadBroker, a_port, 100000, 64 * 1024, clientName);

      }

      FetchRequest req = new FetchRequestBuilder().clientId(clientName).addFetch(a_topic, a_partition, readOffset, 100000).build();

      FetchResponse fetchResponse = consumer.fetch(req);

      if (fetchResponse.hasError()) {
        numErrors++;
        // Something went wrong!
        short code = fetchResponse.errorCode(a_topic, a_partition);
        System.out.println("Error fetching data from the Broker:" + leadBroker + " Reason: " + code);

        if (numErrors > 5)
          break;
        if (code == ErrorMapping.OffsetOutOfRangeCode()) {
          // We asked for an invalid offset. For simple case ask for
          // the last element to reset
          readOffset = getLastOffset(consumer, a_topic, a_partition, kafka.api.OffsetRequest.LatestTime(), clientName);
          continue;
        }

        consumer.close();
        consumer = null;
        leadBroker = findNewLeader(leadBroker, a_topic, a_partition, a_port);
        continue;
      }

      numErrors = 0;

      long numRead = 0;
      for (MessageAndOffset messageAndOffset : fetchResponse.messageSet(a_topic, a_partition)) {
        long currentOffset = messageAndOffset.offset();
        if (currentOffset < readOffset) {
          System.out.println("Found an old offset: " + currentOffset + " Expecting: " + readOffset);
          continue;
        }

        readOffset = messageAndOffset.nextOffset();
        ByteBuffer payload = messageAndOffset.message().payload();

        byte[] bytes = new byte[payload.limit()];
        payload.get(bytes);
        System.out.println(String.valueOf(messageAndOffset.offset()) + ": " + new String(bytes, "UTF-8"));
        numRead++;
        a_maxReads--;
      }

      if (numRead == 0) {
        try {
          Thread.sleep(1000);
        } catch (InterruptedException ie) {
        }
      }
    }

    if (consumer != null)
      consumer.close();
  }

  public static long getLastOffset(SimpleConsumer consumer, String topic, int partition, long whichTime, String clientName) {

    TopicAndPartition topicAndPartition = new TopicAndPartition(topic, partition);

    Map<TopicAndPartition, PartitionOffsetRequestInfo> requestInfo = new HashMap<TopicAndPartition, PartitionOffsetRequestInfo>();

    requestInfo.put(topicAndPartition, new PartitionOffsetRequestInfo(whichTime, 1));

    kafka.javaapi.OffsetRequest request = new kafka.javaapi.OffsetRequest(requestInfo, kafka.api.OffsetRequest.CurrentVersion(), clientName);

    OffsetResponse response = consumer.getOffsetsBefore(request);

    if (response.hasError()) {
      System.out.println("Error fetching data Offset Data the Broker. Reason: " + response.errorCode(topic, partition));
      return 0;
    }
    long[] offsets = response.offsets(topic, partition);
    return offsets[0];
  }

  private String findNewLeader(String a_oldLeader, String a_topic, int a_partition, int a_port) throws Exception {

    for (int i = 0; i < 3; i++) {
      boolean goToSleep = false;
      PartitionMetadata metadata = findLeader(m_replicaBrokers, a_port, a_topic, a_partition);
      if (metadata == null) {
        goToSleep = true;
      } else if (metadata.leader() == null) {
        goToSleep = true;
      } else if (a_oldLeader.equalsIgnoreCase(metadata.leader().host()) && i == 0) {
        // first time through if the leader hasn't changed give
        // ZooKeeper a second to recover
        // second time, assume the broker did recover before failover,
        // or it was a non-Broker issue

        goToSleep = true;
      } else {
        return metadata.leader().host();
      }

      if (goToSleep) {
           Thread.sleep(1000);
      }
    }
    System.out.println("Unable to find new leader after Broker failure. Exiting");
    throw new Exception("Unable to find new leader after Broker failure. Exiting");
  }

  private PartitionMetadata findLeader(List<String> a_seedBrokers, int a_port, String a_topic, int a_partition) {
    PartitionMetadata returnMetaData = null;

    loop:
    for (String seed : a_seedBrokers) {
      SimpleConsumer consumer = null;

      try {
        consumer = new SimpleConsumer(seed, a_port, 100000, 64 * 1024, "leaderLookup");
        List<String> topics = Collections.singletonList(a_topic);
        TopicMetadataRequest req = new TopicMetadataRequest(topics);
        kafka.javaapi.TopicMetadataResponse resp = consumer.send(req);
        List<TopicMetadata> metaData = resp.topicsMetadata();

        for (TopicMetadata item : metaData) {
          for (PartitionMetadata part : item.partitionsMetadata()) {
            if (part.partitionId() == a_partition) {
              returnMetaData = part;
               break loop;
            }
          }
        }
      } catch (Exception e) {
        System.out.println("Error communicating with Broker [" + seed + "] to find Leader for [" + a_topic + ", " + a_partition + "] Reason: " + e);
      } finally {
        if (consumer != null)
          consumer.close();
      }
    }

    if (returnMetaData != null) {
      m_replicaBrokers.clear();
      for (BrokerEndPoint replica : returnMetaData.replicas()) {
        m_replicaBrokers.add(replica.host());
      }
    }
    return returnMetaData;
  }
}

Views: 569