Storm计算网站PV和UV(实现可靠处理)

需求分析

编写Storm拓扑实现可靠计算网站当日PV和UV

重点:

  1. 去重计算模式
  2. 实现可靠处理

电商常用指标之PV、UV、VV、独立IP

  • PV(访问量):Page View, 即页面浏览量或点击量,用户每次访问即被计算一次。
  • UV(独立访客):Unique Visitor, 访问您网站的一台电脑客户端为一个访客。00:00-24:00内相同的客户端只会被计算一次。
  • VV即Visit View,访客访问的次数,用以记录所有访客一天内访问量多少次网站。
  • IP(独立IP):指独立IP数。00:00-24:00内相同IP地址之被计算一次。

如果我们要统计一个网站当日有多少实际用户访问,使用UV也就是独立用户作为统计量有什么好处?它比独立IP的统计更加准确吗?

比如你早上8点使用PC访问了某商城的两个页面,下午2点又使用手机访问了这个商城网站的3个页面, 那么对于这个商城网站, 当日的PV、UV、VV、IP各项指标该如何计算呢?

PV为5, PV指浏览量,因此PV指等于上午浏览的2个页面和下午浏览的3个页面之和;

UV为1, UV指独立访客数,因此一天内同一访客的多次访问只计为1个UV;

VV为2, VV指访客的访问次数,上午和下午分别有一次访问行为,因此VV为2

IP为2, IP为独立IP数,由于PC和手机访问时的IP不同,因此独立IP数为2.

可见IP是一个反映网络虚拟地址对象的概念,UV是一个反映实际使用者的概念,每个UV相对于每个IP更加准确地对应一个实际的浏览者。

方案设计

综上所述:使用UV作为统计量,可以更加准确的了解单位时间内实际上有多少个访问者来到了相应的页面。

用Cookie携带的SessionID分析UV值:当客户端第一次访问某个网站服务器的时候,网站服务器会给这个客户端的电脑发出一个Cookie,通常放在这个客户端电脑的C盘当中。在这个Cookie中会分配一个独一无二的编号,这其中会记录一些访问服务器的信息,如访问时间,访问了哪些页面等等。当你下次再访问这个服务器的时候,服务器就可以直接从你的电脑中找到上一次放进去的Cookie文件,并且对其进行一些更新,但那个独一无二的编号是不会变的。

如果把sessionID 放入Set集合实现自动去重,再通过Set.size() 获得UV。该方案在单线程和单JVM下没有问题。高并发情况的多线程情况不适用。

推荐的方案:编写Storm拓扑, Spout获取外部服务器日志然后再通过shuffleGrouping发给多个下级bolt进行数据清洗, 清洗后的数据通过fieldGrouping 进行多线程局部汇总得到,下级blot进行单线程保存sessionID和count数到Map,下一级blot3进行Map遍历,可以得到:Pv、UV、访问深度(每个session_id 的浏览数)

file

创建IDEA项目(基于Maven)

pom.xml

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <groupId>cn.delucia</groupId>
    <artifactId>uv-with-reliability</artifactId>
    <version>1.0-SNAPSHOT</version>

    <dependencies>
        <dependency>
            <groupId>org.apache.storm</groupId>
            <artifactId>storm-core</artifactId>
            <version>2.1.0</version>
            <scope>compile</scope>
        </dependency>
    </dependencies>

</project>

项目结构

file

代码编写

创建工具类cn.delucia.storm.utils.SocketClientUtil

package cn.delucia.storm.utils;

import java.io.OutputStream;
import java.lang.management.ManagementFactory;
import java.net.InetAddress;
import java.net.Socket;
import java.net.UnknownHostException;

/**
 * @Author: deLucia
 * @Date: 2021/9/29
 * @Version: 1.0
 * @Description: 工具类
 */
public class SocketClientUtil {
    /**
     * 是否启用send方法向Socket Server发送信息
     */
    public static boolean ENABLE = true;
    public static final String HOST = "localhost";
    public static final int PORT = 9876;

    private SocketClientUtil() {
    }

    private static String getHostname() {
        try {
            return InetAddress.getLocalHost().getHostName();
        } catch (UnknownHostException e) {
            e.printStackTrace();
        }
        return null;
    }

    /**
     * 返回进程pid
     */
    private static String getPID() {
        String info = ManagementFactory.getRuntimeMXBean().getName();
        return info.split("@")[0];
    }

    /**
     * 线程信息
     */
    private static String getTID() {
        return Thread.currentThread().getName();
    }

    //对象信息
    private static String getOID(Object obj) {
        String cname = obj.getClass().getSimpleName();
        int hash = obj.hashCode();
        return cname + "@" + hash;
    }

    public static String info(Object obj, String msg) {
        return getHostname() + ", " + getPID() + ", " + getTID() + ", " + getOID(obj) + ", " + msg;
    }

    /**
     * 向远端发送sock消息
     * 远端需要先开启nc进行监听:
     * windows: nc -l localhost -L -p 9876
     * linux: nc -l localhost -k -p 9876
     */
    public static void send(Object obj, String msg) {
        if (!ENABLE) {
            return;
        }
        String info = info(obj, msg);
        try (Socket sock = new Socket(HOST, PORT); OutputStream os = sock.getOutputStream()) {
            os.write((info + "\r\n").getBytes());
            os.flush();
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

创建类cn.delucia.storm.log.LogTopology

package cn.delucia.storm.log;

import cn.delucia.storm.utils.SocketClientUtil;
import com.google.common.collect.ImmutableList;
import org.apache.storm.Config;
import org.apache.storm.LocalCluster;
import org.apache.storm.StormSubmitter;
import org.apache.storm.topology.TopologyBuilder;
import org.apache.storm.tuple.Fields;

/**
 * @author Delucia
 */
public class LogTopology {

    public static final String HOST = "hadoop000";
    public static final Integer ZK_SERVER_PORT = 2181;
    public static final String SPOUT1_ID = LogMockSpout.class.getSimpleName();
    public static final String BOLT1_ID = LogCheckBolt.class.getSimpleName();
    public static final String BOLT2_ID = LogCountBolt.class.getSimpleName();
    public static final String BOLT3_ID = LogGlobalSumBolt.class.getSimpleName();

    public static void main(String[] args) throws Exception {

        // 构造拓扑 DAG(有向无环图)
        TopologyBuilder builder = new TopologyBuilder();
        builder.setSpout(SPOUT1_ID, new LogMockSpout(), 1);
        builder.setBolt(BOLT1_ID, new LogCheckBolt(), 4).shuffleGrouping(SPOUT1_ID);
        builder.setBolt(BOLT2_ID, new LogCountBolt(), 4).fieldsGrouping(BOLT1_ID,
                new Fields("date", "sid"));
        // 单线程全局汇总
        builder.setBolt(BOLT3_ID, new LogGlobalSumBolt(), 1).globalGrouping(BOLT2_ID);
        // Storm相关配置
        Config conf = new Config();
        // TOPOLOGY_DEBUG: ON 日志会详细记录所有发射的数据
        conf.setDebug(false);
        // Executor number for event loggers
        conf.setNumEventLoggers(0);
        // Executor number for ackers
        conf.setNumAckers(1);
        // number of workers
        conf.setNumWorkers(1);
        // transfer queue between workers
        conf.put(Config.TOPOLOGY_TRANSFER_BUFFER_SIZE, 64);
        // reveicer queue for each executor
        conf.put(Config.TOPOLOGY_EXECUTOR_RECEIVE_BUFFER_SIZE, 16384);

        if (args.length > 0 ) {
            SocketClientUtil.ENABLE = false;
            conf.put(Config.STORM_ZOOKEEPER_SERVERS, ImmutableList.of(HOST));
            conf.put(Config.STORM_ZOOKEEPER_PORT, ZK_SERVER_PORT);
            conf.put(Config.NIMBUS_SEEDS, ImmutableList.of(HOST));
            conf.put(Config.NIMBUS_THRIFT_PORT, 16627);
            String jarLocation = ".\\target\\parallelisim-1.0-SNAPSHOT.jar";
            System.setProperty("storm.jar", jarLocation);

            // 使用 StormSubmitter 提交本地构建的jar包
            StormSubmitter.submitTopology(args[0], conf, builder.createTopology());

        } else {
            LocalCluster cluster = new LocalCluster();
            cluster.submitTopology(LogTopology.class.getSimpleName(), conf, builder.createTopology());
        }
    }
}

创建类cn.delucia.storm.log.LogMockSpout

package cn.delucia.storm.log;

import cn.delucia.storm.utils.SocketClientUtil;
import org.apache.storm.spout.SpoutOutputCollector;
import org.apache.storm.task.TopologyContext;
import org.apache.storm.topology.OutputFieldsDeclarer;
import org.apache.storm.topology.base.BaseRichSpout;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Values;
import org.apache.storm.utils.Utils;

import java.util.Map;
import java.util.Random;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;

/**
 * @Author: deLucia
 * @Date: 2021/9/28
 * @Version: 1.0
 * @Description: 统计当日日志中的PV和UV指标
 * PV: Page View 每次访问技术一次
 * UV: Unique View 相同用户无论访问多少次都只计数1次
 * 如何判断是否为相同用户:
 * - IP 用户可能更换设备,或者使用VPN工具, 因此根据IP的统计可能不准确
 * - Cookie(SessionID) 一般使用cookie判断
 */
public class LogMockSpout extends BaseRichSpout {
    /**
     * Mock Data
     * 注意这里伪造的日期需要有当天的日期,否则最终统计的PV和UV结果都为 0
     * 假设今天是 2021年10月2日。
     * 2022年的日期因为晚于今天, 因此是无效数据, 需要重试3次
     */
    private final String[] dates = {"2021-10-02 08:40:50", "2022-10-02 18:40:50", "2022-10-02 19:40:50", "2021-10-02 20:40:50"};
    private final String[] sids = {"user100", "user101", "user102", "user103"};
    private final String[] urls = {"https://cn.delucia/a.html", "https://cn.delucia/b.html", "https://cn.delucia/c.html", "http://cn.delucia/d.html"};
    /**
     * Spout发射的tuple数量
     */
    private static final Integer EMMIT_TUPLE_NUMBER = 30;
    /**
     * tuple处理失败后进行重试的次数
     */
    private Integer retryTimes = 0;
    /**
     * 允许的最大重试次数
     */
    private static final Integer MAX_RETRY_TIMES = 3;
    private Integer count = 0;
    /**
     * 缓存所有发送的消息, 消息发射后移除
     */
    private ConcurrentHashMap<Object, String> cachedMessages;
    /**
     * 用于发射tuple到下游
     */
    private transient SpoutOutputCollector collector;
    /**
     * 用于生产随机数
     */
    private Random rand;
    /**
     * 自定义线程池 - 用于延迟发送需要重试的数据
     */
    private transient ThreadPoolExecutor executor;

    /**
     * 每个线程初始化时调用
     *
     * @param conf
     * @param context
     * @param collector
     */
    @Override
    public void open(Map conf, TopologyContext context, SpoutOutputCollector collector) {
        this.collector = collector;
        this.cachedMessages = new ConcurrentHashMap<>();
        this.rand = new Random();

        // 自定义线程池 - 用于执行重试任务
        int corePoolSize = 2;
        int maximumPoolSize = 4;
        long keepAliveTime = 10;
        TimeUnit unit = TimeUnit.SECONDS;
        BlockingQueue<Runnable> workQueue = new ArrayBlockingQueue<>(3);
        ThreadFactory threadFactory = new NameThreadFactory();
        RejectedExecutionHandler handler = new IgnorePolicy();

        executor = new ThreadPoolExecutor(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, threadFactory, handler);
        executor.prestartAllCoreThreads(); // 预启动所有核心线程
    }

    @Override
    public void nextTuple() {
        if (count >= EMMIT_TUPLE_NUMBER) {
            return;
        }
        int k = rand.nextInt(4);
        int l = rand.nextInt(4);
        int m = rand.nextInt(4);
        String line = dates[k] + "\t" + sids[l] + "\t" + urls[m];
        String msg = ">> Emit Tuple[\"line\"]: " + line;
        System.err.println(msg);
        SocketClientUtil.send(this, msg);
        Long time = System.currentTimeMillis();
        this.collector.emit(new Values(line), time);
        cachedMessages.put(time, line);
        Utils.sleep(500);
        count++;
    }

    @Override
    public void declareOutputFields(OutputFieldsDeclarer declarer) {
        declarer.declare(new Fields("line"));
    }

    @Override
    public void ack(Object msgId) {
        cachedMessages.remove(msgId);
    }

    @Override
    public void fail(Object msgId) {
        cachedMessages.forEach((id, message) -> {
            if (msgId.equals(id)) {
                cachedMessages.put(id, message);
                retryTimes += 1;
                // 重试
                if (retryTimes <= MAX_RETRY_TIMES) {
                    // 同步改异步
                    // collector.emit(new Values(message), msgId);
                    RetryTask retryTask = new RetryTask(collector, message, msgId, retryTimes);
                    executor.execute(retryTask);
                } else {
                    // 超过最大重试次数, 放弃
                    retryTimes = 0;
                    cachedMessages.remove(msgId);
                }
            }
        });
    }

    class NameThreadFactory implements ThreadFactory {

        private final AtomicInteger mThreadNum = new AtomicInteger(1);

        @Override
        public Thread newThread(Runnable r) {
            Thread t = new Thread(r, "my-thread-" + mThreadNum.getAndIncrement());
            // String msg = "Start Thread[" + t.getName() + "] will retry in 30 secs, max retry times is " + MAX_RETRY_TIMES + ".";
            // System.err.println(msg);
            // SocketClientUtil.send(this, msg);
            return t;
        }
    }

    static class RetryTask implements Runnable {
        private final SpoutOutputCollector collector;
        private final String messageBody;
        private final Object messageId;
        private final Integer retryTimes;
        private final Random rand;

        public RetryTask(SpoutOutputCollector collector, String name, Object msgId, Integer retryTimes) {
            this.collector = collector;
            this.messageBody = name;
            this.messageId = msgId;
            this.retryTimes = retryTimes;
            this.rand = new Random();
        }

        @Override
        public void run() {
            // 30秒内进行重试
            Utils.sleep(rand.nextInt(30000));
            collector.emit(new Values(messageBody), messageId);
            String msg = ">> Replay: " + retryTimes + " times: " + messageBody;
            System.err.println(msg);
            SocketClientUtil.send(this, msg);
        }

        @Override
        public String toString() {
            return "Task: " + messageBody;
        }
    }

    public static class IgnorePolicy implements RejectedExecutionHandler {

        @Override
        public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
            log(r, e);
        }

        private void log(Runnable r, ThreadPoolExecutor e) {
            String msg = r.toString() + " rejected";
            System.err.println(msg);
            SocketClientUtil.send(this, msg);
        }
    }
}

创建类cn.delucia.storm.log.LogCheckBolt

package cn.delucia.storm.log;

import org.apache.storm.task.OutputCollector;
import org.apache.storm.task.TopologyContext;
import org.apache.storm.topology.OutputFieldsDeclarer;
import org.apache.storm.topology.base.BaseRichBolt;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.Values;

import java.text.ParseException;
import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.Map;

/**
 * @Author: deLucia
 * @Date: 2021/9/28
 * @Version: 1.0
 * @Description: 检查日志的格式是否符合要求
 */
public class LogCheckBolt extends BaseRichBolt {
    private static final long serialVersionUID = 1L;

    OutputCollector collector;

    @Override
    public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) {
        this.collector = collector;
    }

    @Override
    public void execute(Tuple input) {

        String line = input.getStringByField("line");
        if (line == null || line.isEmpty()) {
            this.collector.fail(input);
            return;
        }
        String date = line.split("\t")[0];
        String sid = line.split("\t")[1];

        SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd");
        try {
            Date d = sdf.parse(date);
            if (d.after(new Date())) {
                this.collector.fail(input);
                return;
            }
            date = sdf.format(d);
            this.collector.emit(input, new Values(date, sid));
            this.collector.ack(input);
        } catch (ParseException e) {
            e.printStackTrace();
            // 日期格式不正确
            this.collector.fail(input);
        }
    }

    @Override
    public void declareOutputFields(OutputFieldsDeclarer declarer) {
        declarer.declare(new Fields("date", "sid"));
    }
}

创建类cn.delucia.storm.log.LogCountBolt

package cn.delucia.storm.log;

import org.apache.storm.task.OutputCollector;
import org.apache.storm.task.TopologyContext;
import org.apache.storm.topology.OutputFieldsDeclarer;
import org.apache.storm.topology.base.BaseRichBolt;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.Values;

import java.util.HashMap;
import java.util.Map;

/**
 * @Author: deLucia
 * @Date: 2021/9/28
 * @Version: 1.0
 * @Description: 对相同的sessionID进行聚合计数
 */
public class LogCountBolt extends BaseRichBolt {

    OutputCollector collector;
    // <session_id>, <count>
    Map<String, Long> map = new HashMap<>();

    @Override
    public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) {
        this.collector = collector;
    }

    @Override
    public void execute(Tuple input) {

        String date = input.getStringByField("date");
        String ip = input.getStringByField("sid");
        String key = date + "_" + ip;
        Long count = 0L;
        try {
            count = map.get(key);
            if (count == null) {
                count = 0L;
            }
            count++;
            map.put(key, count);
            this.collector.emit(input, new Values(date + "_" + ip, count));
            this.collector.ack(input);
        } catch (Exception e) {
            e.printStackTrace();
            System.err.println(LogCountBolt.class.getSimpleName() + " failed, date: " + date + ", sid: " + ip + ", count: " + count);
            this.collector.fail(input);
        }
    }

    @Override
    public void declareOutputFields(OutputFieldsDeclarer declarer) {
        declarer.declare(new Fields("date_sid", "count"));
    }
}

创建类cn.delucia.storm.log.LogGlobalSumBolt

package cn.delucia.storm.log;

import cn.delucia.storm.utils.SocketClientUtil;
import org.apache.storm.task.OutputCollector;
import org.apache.storm.task.TopologyContext;
import org.apache.storm.topology.OutputFieldsDeclarer;
import org.apache.storm.topology.base.BaseRichBolt;
import org.apache.storm.tuple.Tuple;

import java.text.ParseException;
import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.HashMap;
import java.util.Map;

/**
 * @Author: deLucia
 * @Date: 2021/9/28
 * @Version: 1.0
 * @Description: 全局聚合 - 只统计当日的 pv 和 uv
 */
public class LogGlobalSumBolt extends BaseRichBolt {
    SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd");
    private static final long serialVersionUID = 1L;
    transient OutputCollector collector;
    Map<String, Long> map = new HashMap<>();

    @Override
    public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) {
        this.collector = collector;
    }

    @Override
    public void execute(Tuple input) {
        int pv = 0;
        int uv = 0;

        String dateIp = input.getStringByField("date_sid");
        Long count = input.getLongByField("count");
        String currDate = sdf.format(new Date());
        collector.ack(input);
        try {
            Date accessDate = sdf.parse(dateIp.split("_")[0]);
            // 访问日期不能晚于今天
            map.put(dateIp, count);
            String msg = ">> date_sid: " + dateIp + ", pv: " + count;
            System.err.println(msg);
            SocketClientUtil.send(this, msg);

            // 只计算当日的pv和uv
            if (!dateIp.split("_")[0].startsWith(currDate)) {
                return;
            }
            if (!map.isEmpty()) {
                for (Map.Entry<String, Long> e : map.entrySet()) {
                    uv++;
                    pv += e.getValue();
                }
            }
            String result = ">> " + currDate + " pv: " + pv + ", uv: " + uv;
            System.err.println(result);
            SocketClientUtil.send(this, result);
        } catch (ParseException e1) {
            e1.printStackTrace();
        } finally {
            // 到此整棵tuple树处理完成
            collector.ack(input);
        }

    }

    @Override
    public void declareOutputFields(OutputFieldsDeclarer declarer) {

    }

}

本地运行

在本机安装netcat,使用一下命令开启Socket Server:

>nc -l localhost -L -p 9876

然后本地运行项目(不要带参数),等待Socket Server接收到输出结果,如图所示:

file

本地提交打好的jar包提交到远程集群运行

首先在pom.xml中将依赖storm-core的scope改为provided
再使用maven package打包,生成的文件在项目根目录下的target文件夹内
修改LogTopology.java中的代码,指定jar包的本地提交路径:

String jarLocation = ".\\target\\parallelisim-1.0-SNAPSHOT.jar";
System.setProperty("storm.jar", jarLocation);

然后直接运行项目,即可通过Thrift协议协议将jar包远程提交到集群中的nimbus节点运行。
最后可以在Storm UI界面查看结果。

Views: 351

Storm拓扑之Stream Grouping

Stream的分组策略

Stream Grouping - 定义了一个流在Bolt任务间该如何被切分,谁来处理哪些数据流,按照什么规则来分配.

随机分组

Shuffle Grouping- 随机分组, 随机派发stream里面的tuple,保证每个bolt接收到的tuple数目大致相同。

字段分组

Fields grouping - 根据指定字段的值进行分组。比如说,一个数据流根据'word'字段进行分组,所有具有相同的'word'字段值的tuple会路由到同一个bolt的task中。

广播分组

All grouping- 全复制分组(广播分组), 将所有的tuple复制后分发给所有的bolt task。每个订阅数据流的task都会接收到tuple的拷贝。

全局分组

Global Grouping- 全局分组:这种分组方式将所有的tuple路由到唯一一个task上。Storm按照最小的taskID来选取接收数据的task。注意!!当使用全局分组方式时,设置bolt的task并发度是没有意义的,因为所有tuple都转发到同一个task上了。使用全局分组的时候需要注意,因为所有tuple都转发到一个JVM实例上,可能会引起Storm集群中某个JVM或者服务器出现性能瓶颈或奔溃。

无分组

None Grouping - 不分组,这个分组的意思是说stream不关心到底谁会收到它的tuple。目前这种分组和Shuffle grouping是一样的效果, 有一点不同的是storm会把这个bolt放到这个bolt的订阅者同一个线程里面去执行。

直接分组

Direct Grouping — 直接分组, 这是一种比较特别的分组方法,用这种分组意味着消息的发送者指定由消息接收者的哪个task处理这个消息。 只有被声明为Direct Stream的消息流可以声明这种分组方法。而且这种消息tuple必须使用emitDirect方法来发射。消息处理者可以通过TopologyContext来获取处理它的消息的task的id (OutputCollector.emit方法也会返回task的id)。

本地或随机分组

Local or shuffle grouping - 本地或随机分组:和随机分组类似,但是,会将tuple分发给同一个worker内的bolt task(如果worker内有接收数据的bolt task)。其他情况下,采用随机分组的方式。取决于topology的并发度,本地或随机分组可以减少网络传输,从而提高topology性能。

作业:

自己探索并实现其中任意一种分组方式,并添加说明性注释和运行截图

代码下载:
1494-2021-09-18-storm_Chap2Ex

Views: 274

Storm 累加拓扑示例

创建Spout发送递增数字数列

private static class NumSpout extends BaseRichSpout {
    private SpoutOutputCollector collector;
    private int num = 1;

    @Override
    public void open(Map<String, Object> conf, TopologyContext context, SpoutOutputCollector collector) {

        this.collector = collector;
    }

    @Override
    public void nextTuple() {

        collector.emit(new Values(num));
        System.err.println("+ " + num);
        Utils.sleep(1000);
        num++;
    }

    @Override
    public void declareOutputFields(OutputFieldsDeclarer declarer) {
        declarer.declare(new Fields("num"));
    }
}

创建Bolt负责计算累加结果

    private static class SumBolt extends BaseRichBolt {

        // private OutputCollector collector;

        private int sum = 0;
        @Override
        public void prepare(Map<String, Object> topoConf, TopologyContext context, OutputCollector collector) {

            // this.collector = collector;
        }

        @Override
        public void execute(Tuple input) {
            //int number = input.getInteger(0);
            int number = input.getIntegerByField("num");
            sum += number;
            // collector.emit(new Values(sum));
            System.err.println("= " + sum);
        }

        @Override
        public void declareOutputFields(OutputFieldsDeclarer declarer) {

            // declarer.declare(new Fields("sum"));
        }
    }

本地运行

本地运行完整代码

public class MySumTopology {

    public static void main throws Exception (String[] args) {

        Config config = new Config();
        config.setNumWorkers(1);
        config.setDebug(true);

        TopologyBuilder topologyBuilder = new TopologyBuilder();
        topologyBuilder.setSpout("num-spout", new NumSpout());
        topologyBuilder.setBolt("sum-bolt", new SumBolt()).shuffleGrouping("num-spout");

        LocalCluster localCluster = new LocalCluster();
        localCluster.submitTopology("sum-topo", config, topologyBuilder.createTopology());
        // 运行1分钟停止
        Utils.sleep(60000);
        localCluster.shutdown();
    }
}

保证数据可靠处理

  1. Spout在使用nextTuple()方法发送数据时需要传入消息ID
    @Override
    public void nextTuple() {
        collector.emit(new Values(num), num);
        ...
    }
  1. Bolt中execute()方法中标记tuple是否处理成功
  • 处理成功 collector.ack(input)

  • 处理失败collector.fail(input)

    注意: ackfail方法需要锚定到发射过来的tuple上.

    private static class SumBolt extends BaseRichBolt {

        private OutputCollector collector;
        private Random rand;

        private int sum = 0;

        @Override
        public void execute(Tuple input) {

            if (rand.nextInt(10) < 2) { 
                collector.fail(input); // 20%几率 处理失败 
                System.err.println("!搞错了");
                return;
            }

            int number = input.getIntegerByField("num");
            sum += number;
            // collector.emit(new Values(sum));
            System.err.println("= " + sum);
            collector.ack(input); // 处理成功: Acknowledge
        }
    ...
  1. Spout中对处理失败的元组触发回调

    这里把处理失败的数字重新发射出去

    @Override
    public void fail(Object msgId) {
        int oldNum = num - 1;
        collector.emit(new Values(oldNum), msgId);
        System.out.println("重来 + " + oldNum);
        Utils.sleep(1000);
    }

远程运行模式

如果是需要提交到集群, 需要使用StormSubmitter替换new LocalCluster()来调用submitTopology(..)方法

// local mode运行 - LocalCluster
LocalCluster cluster = new LocalCluster();
cluster.submitTopology("sum-topo", config, builder.createTopology());

// remote mode - run on a production cluster
// StormSubmitter.submitTopology("sum-topo",config,builder.createTopology());

然后打成jar包上传到集群运行

Views: 222

Kafka伪集群环境搭建

创建Zookeeper集群(3个)

前提是已经装好Java JDK8+并配置好环境变量。

建议Kafka集群使用专有的Zookeeper集群进行协调管理。

也可以使用Kafka内置的bin/zookeeper命令启动集群, 默认配置是config/zookeeper.properties

创建3个zk配置文件

[hadoop@hadoop000 config] vi zookeeper-1(2|3).properties

修改配置文件内容如下

# The number of milliseconds of each tick
tickTime=2000
# The number of ticks that the initial 
# synchronization phase can take
initLimit=10
# The number of ticks that can pass between 
# sending a request and getting an acknowledgement
syncLimit=5
# the directory where the snapshot is stored.
dataDir=/home/hadoop/tmp/kafka-zk-1(2|3)/data
dataLogDir=/home/hadoop/tmp/kafka-zk-1(2|3)/log
# the port at which the clients will connect
clientPort=2181(2|3)
# 将所有4字命令添加到白名单中
4lw.commands.whitelist=*
# 配置zookeeper服务器的 [主机名或主机ID]:[同步端口(唯一)]:[选举端口(唯一)]
server.1=localhost:2891:3891
server.2=localhost:2892:3892
server.3=localhost:2893:3893

根据配置创建对应的dataDir以及dataLogDir,并在dataDir下创建myid文件。

[hadoop@hadoop000 tmp]echo 1 > kafka-zk-1/data/myid
[hadoop@hadoop000 tmp]echo 2 > kafka-zk-2/data/myid
[hadoop@hadoop000 tmp]echo 3 > kafka-zk-3/data/myid

[hadoop@hadoop000 config] zookeeper-server-start.sh -daemon zookeeper-1.properties
[hadoop@hadoop000 config] zookeeper-server-start.sh -daemon zookeeper-2.properties
[hadoop@hadoop000 config]$ zookeeper-server-start.sh -daemon zookeeper-3.properties

修改Kafka集群配置

vi server-1(2|3).properties

# The id of the broker. 
broker.id=1(2|3)
# The port of the broker. 
# port=9092(3|4)
listeners=PLAINTEXT://hadoop000:9092(3|4)
# Log file directories - A comma separated list
log.dirs=/home/hadoop/tmp/kafka-1(2|3)/log
# Zookeeper connection string - A comma separated list
zookeeper.connect=localhost:2181,localhost:2182,localhost:2183/kafka

启动Kafka代理服务(3个)

[hadoop@hadoop000 config]$ kafka-server-start.sh -daemon server-1.properties 
[hadoop@hadoop000 config]$ kafka-server-start.sh -daemon server-2.properties 
[hadoop@hadoop000 config]$ kafka-server-start.sh -daemon server-3.properties 

创建消息主题(2分区,2副本)

由于我们有3台Kafka服务器,因此可以创建具有多分区以及多副本的主题

$ bin/kafka-topics.sh --bootstrap-server hadoop000:9092,hadoop000:9093,hadoop000:9094 --create --topic my-topic --partitions 2 --replication-factor 2

也可以使用--zookeeper选项进行连接, 如下所示:

$ kafka-topics.sh --zookeeper localhost:2181,localhost:2182,localhost:2183/kafka --create --topic my-topic  --partitions 2 --replication-factor 2

但是--zookeeper这个选项在较新版本中已经废弃, 建议使用--bootstrap来代替, 有更好的安全机制

列出主题详情

$ bin/kafka-topics.sh --bootstrap-server hadoop000:9092,hadoop000:9093,hadoop000:9094 --describe --topic my-topic

Topic: my-topic PartitionCount: 2   ReplicationFactor: 2    Configs: 
Topic: my-topic Partition: 0    Leader: 2   Replicas: 2,3   Isr: 2,3
Topic: my-topic Partition: 1    Leader: 3   Replicas: 3,1   Isr: 3,1

创建消费者

可以同时监听多台Kafka服务器组成的集群

  bin/kafka-console-consumer.sh --bootstrap-server hadoop000:9092,hadoop000:9093,hadoop000:9094 --topic my-topic --from-beginning

创建生产者

  • 代理服务器列表可以指定多台Kafka服务器组成的集群

  • 按顺序发出一些消息

    kafka-console-producer.sh --broker-list hadoop000:9092,hadoop000:9093,hadoop000:9094 --topic my-topic
    0
    1
    2
    3
    4
    5

    消费者查看消息

[hadoop@hadoop000 ~]$ kafka-console-consumer.sh --bootstrap-server hadoop000:9092,hadoop000:9093,hadoop000:9094 --topic my-topic --from-beginning
4
5
0
1
2
3

小问题:消费者客户端接收的消息顺序为什么这样?

因为只有消费者接受的消息都是来自单个分区的时候才能保证消息有序.

总结

以上安装方式虽然使用了三个zookeeper服务器和三个kafka broker,但是还是运行在一台机器上,因此只能算得上伪分布式,真正的分布式需要将这些服务分布在多台机器上的。

Views: 383

Kafka集群部署的讨论

只有单台机器构成的 Kafka 伪集群只能用于日常测试之用,根本无法满足实际的线上生产需求。而真正的线上环境需要仔细地考量各种因素,结合自身的业务需求而制定。下面我就分别从操作系统、磁盘、磁盘容量和带宽等方面来讨论一下。

操作系统

首先我们先看看要把 Kafka 安装到什么操作系统上。

目前常见的操作系统有 3 种:

  1. Linux
  2. Windows
  3. macOS。

如果考虑操作系统与 Kafka 的适配性,Linux 系统显然要比其他两个特别是 Windows 系统更加适合部署 Kafka。主要是在下面这三个方面上,Linux 的表现更胜一筹。

  • I/O 模型的使用
    • Kafka 客户端底层使用了 Java 的 selector,selector 在 Linux 上的I/O模型的实现机制和Windows是不同的。部署在 Linux 上比起部署在Windows上能够获得更高效的 I/O 性能。
  • 数据网络传输效率
    • 在 Linux 部署 Kafka 能够享受到零拷贝技术所带来的快速数据传输特性。
  • 社区支持度
    • 相比Windows系统,社区会优先解决Linux系统上发现的Bug。

总结:Windows 平台上部署 Kafka 只适合于个人测试或用于功能验证,千万不要应用于生产环境。

磁盘

磁盘对 Kafka 性能的影响最重要。在对 Kafka 集群进行磁盘规划时经常面对的问题是,我应该选择普通的机械磁盘还是固态硬盘?前者成本低且容量大,但易损坏;后者性能优势大,不过单价高。

Kafka 大量使用磁盘不假,可它使用的方式多是顺序读写操作,一定程度上规避了机械磁盘最大的劣势,即随机读写操作慢。从这一点上来说,使用 SSD 似乎并没有太大的性能优势,毕竟从性价比上来说,机械磁盘物美价廉,而它因易损坏而造成的可靠性差等缺陷,又由 Kafka 在软件层面提供机制来保证,故使用普通机械磁盘是很划算的。

关于磁盘选择另一个经常讨论的话题就是到底是否应该使用磁盘阵列(RAID)。使用 RAID 的两个主要优势在于:

  • 提供冗余的磁盘存储空间
  • 提供负载均衡

不过就 Kafka 而言,一方面 Kafka 自己实现了冗余机制来提供高可靠性;另一方面通过分区的概念,Kafka 也能在软件层面自行实现负载均衡。

综合以上的考量,建议:

  • 追求性价比的公司可以不搭建 RAID,使用普通磁盘组成存储空间即可。

  • 使用机械磁盘完全能够胜任 Kafka 线上环境。

磁盘容量

Kafka 集群到底需要多大的存储空间?Kafka 需要将消息保存在底层的磁盘上,这些消息默认会被保存一段时间然后自动被删除。虽然这段时间是可以配置的,但你应该如何结合自身业务场景和存储需求来规划 Kafka 集群的存储容量呢?

我举一个简的例子来说明该如何思考这个问题。假设你所在公司有个业务每天需要向 Kafka 集群发送 1 亿条消息,每条消息保存两份以防止数据丢失,另外消息默认保存两周时间。现在假设消息的平均大小是 1KB,那么你能说出你的 Kafka 集群需要为这个业务预留多少磁盘空间吗?

我们来计算一下:

每天, 1 亿条 1KB 大小的消息保存两份且留存两周的时间,那么每天所需空间大小就等于

100000000 * 1KB * 2 / 1000 / 1000 = 200GB

一般情况下 Kafka 集群除了消息数据还有其他类型的数据,比如索引数据等,故我们再为这些数据预留出 10% 的磁盘空间,因此总的存储容量就是

200GB x (1 + 10%) = 220GB

既然要保存两周,那么整体容量即为:

220GB * 14 ≈ 3TB

Kafka 支持数据的压缩,假设压缩比是 0.75,那么最后你需要规划的存储空间就是:

3TB * 0.75 = 2.25TB

总之在规划磁盘容量时你需要考虑下面这几个元素:

  • 新增消息数
  • 消息留存时间
  • 平均消息大小
  • 备份数
  • 是否启用压缩

带宽

对于 Kafka 这种通过网络大量进行数据传输的框架而言,带宽特别容易成为瓶颈。带宽也主要有两种:

  1. 1Gbps 的千兆网络

  2. 10Gbps 的万兆网络

特别是千兆网络应该是一般公司网络的标准配置。其实真正要规划的是所需的 Kafka 服务器的数量。假设你公司的机房环境是千兆网络,即 1Gbps。

现在你有个业务,其业务目标 是在 1 小时内处理 1TB 的业务数据。那么问题来了,你到底需要多少台 Kafka 服务器来完成这个业务呢?

计算一下:

由于带宽是 1Gbps,即每秒处理 1Gb 的数据。

真实环境中每台 Kafka 服务器都是安装在专属的机器上。通常情况下可以按照Kafka 会用到 70% 的带宽资源来计算,因为根据实际使用经验,超过 70% 的阈值就有网络丢包的可能性了,故 70% 的设定是一个比较合理的值,也就是说:单台 Kafka 服务器最多也就能使用大约 700Mb(兆比特) 的带宽资源。

这只是它能使用的最大带宽资源,通常要再额外预留出 2/3 的资源,即:

单台服务器使用带宽 700Mb * 1/3 ≈ 240Mbps

好了,有了 240Mbps,我们就可以计算 1 小时内处理 1TB 数据所需的服务器数量了。根据这个目标,我们每秒需要处理 (1024*1024/60/60)MB*8 = 2330Mb 的数据,除以 240,约等于 10 台服务器。如果消息还需要额外复制两份,那么总的服务器台数还要乘以 3,即 30 台。

用这种方法评估线上环境的服务器台数是比较合理的,而且这个方法能够随着你业务需求的变化而动态调整。

小结

在一开始就应该思考好实际场景下业务所需的集群环境。在考量部署方案时需要通盘考虑,不能仅从单个维度上进行评估。

img

Kafka集群配置

配置并不单单指 Kafka 服务器端的配置,其中既有 Broker 端参数,也有主题(后面我用我们更熟悉的 Topic 表示)级别的参数、JVM 端参数和操作系统级别的参数。接下来将要介绍的所有参数都是那些要修改默认值的参数,因为它们的默认值不适合一般的生产环境。

Broker 端参数

目前 Kafka Broker 提供了近 200 个参数,下面我们按照大的用途类别一组一组地介绍它们,希望可以更有针对性,也更方便你记忆。

  1. 存储信息的重要参数

首先 Broker 是需要配置存储信息的,即 Broker 使用哪些磁盘。那么针对存储信息的重要参数有以下这么几个:

  • log.dirs:这是非常重要的参数,指定了 Broker 需要使用的若干个文件目录路径。要知道这个参数是没有默认值的,这说明什么?这说明它必须由你亲自指定。
  • log.dir:注意这是单数,结尾没有 s,说明它只能表示单个路径,它是补充上一个参数用的。

这两个参数应该怎么设置呢?很简单,你只要设置log.dirs,即第一个参数就好了,不要设置log.dir。而且更重要的是,在线上生产环境中一定要为log.dirs配置多个路径,具体格式是一个 CSV 格式,也就是用逗号分隔的多个路径,比如/home/kafka1,/home/kafka2,/home/kafka3这样。如果有条件的话你最好保证这些目录挂载到不同的物理磁盘上。这样做有两个好处:

  • 提升读写性能:比起单块磁盘,多块物理磁盘同时读写数据有更高的吞吐量。
  • 能够实现故障转移:即 Failover。这是 Kafka 1.1 版本新引入的强大功能。在以前,只要 Kafka Broker 使用的任何一块磁盘挂掉了,整个 Broker 进程都会关闭。但是自 1.1 开始,这种情况被修正了,坏掉的磁盘上的数据会自动地转移到其他正常的磁盘上,而且 Broker 还能正常工作。
  1. ZooKeeper 相关的设置

下面说说与 ZooKeeper 相关的设置。首先 ZooKeeper 是做什么的呢?它是一个分布式协调框架,负责协调管理并保存 Kafka 集群的所有元数据信息,比如集群都有哪些 Broker 在运行、创建了哪些 Topic,每个 Topic 都有多少分区以及这些分区的 Leader 副本都在哪些机器上等信息。

Kafka 与 ZooKeeper 相关的最重要的参数当属zookeeper.connect。这也是一个 CSV 格式的参数,比如我可以指定它的值为zk1:2181,zk2:2181,zk3:2181。2181 是 ZooKeeper 的默认端口。

现在问题来了,如果我让多个 Kafka 集群使用同一套 ZooKeeper 集群,那么这个参数应该怎么设置呢?这时候 chroot 就派上用场了。这个 chroot 是 ZooKeeper 的概念,类似于别名。

如果你有两套 Kafka 集群,假设分别叫它们 kafka1 和 kafka2,那么两套集群的zookeeper.connect参数可以这样指定:zk1:2181,zk2:2181,zk3:2181/kafka1zk1:2181,zk2:2181,zk3:2181/kafka2。切记 chroot 只需要写一次,而且是加到最后的。我经常碰到有人这样指定:zk1:2181/kafka1,zk2:2181/kafka2,zk3:2181/kafka3,这样的格式是不对的。

  1. Broker 连接相关

第三组参数是与 Broker 连接相关的,即客户端程序或其他 Broker 如何与该 Broker 进行通信的设置。有以下三个参数:

  • listeners

    监听器,其实就是告诉外部连接者要通过什么协议访问指定主机名和端口开放的 Kafka 服务。

  • advertised.listeners

    和 listeners 相比多了个 advertised。Advertised 的含义表示宣称的、公布的,就是说这组监听器是 Broker 用于对外发布的。监听器从构成上来说,它是若干个逗号分隔的三元组,每个三元组的格式为<协议名称,主机名,端口号>。这里的协议名称可能是标准的名字,比如 PLAINTEXT 表示明文传输、SSL 表示使用 SSL 或 TLS 加密传输等;也可能是你自己定义的协议名字,比如CONTROLLER: //localhost:9092

    一旦你自己定义了协议名称,你必须还要指定listener.security.protocol.map参数告诉这个协议底层使用了哪种安全协议,比如指定listener.security.protocol.map=CONTROLLER:PLAINTEXT表示CONTROLLER这个自定义协议底层使用明文不加密传输数据。

  • host.name/port

    这两个是过期的参数了。如果主机名和端口号已经通过监听器配置,则可以忽略。

遇到类似host.name这个设置中到底使用 IP 地址还是主机名。统一的建议是:最好全部使用主机名,即 Broker 端和 Client 端应用配置中全部填写主机名。 Broker 源代码中也使用的是主机名,如果你在某些地方使用了 IP 地址进行连接,可能会发生无法连接的问题。

因此需要保证在hosts文件中建立主机名和IP之间的映射

  1. 关于 Topic 管理的

第四组参数是关于 Topic 管理的。我来讲讲下面这三个参数:

  • auto.create.topics.enable:是否允许自动创建 Topic。

    建议最好设置成 false,即不允许自动创建 Topic。在我们的线上环境里面有很多名字稀奇古怪的 Topic,我想大概都是因为该参数被设置成了 true 的缘故。

    你可能有这样的经历,要为名为 test 的 Topic 发送事件,但是不小心拼写错误了,把 test 写成了 tst,之后启动了生产者程序。恭喜你,一个名为 tst 的 Topic 就被自动创建了。

    所以我一直相信好的运维应该防止这种情形的发生,特别是对于那些大公司而言,每个部门被分配的 Topic 应该由运维严格把控,决不能允许自行创建任何 Topic。

  • unclean.leader.election.enable:是否允许 Unclean Leader 选举。

    Kafka每个分区都有多个副本来提供高可用。在这些副本中只能有一个副本对外提供服务,即所谓的 Leader 副本。只有保存数据比较多的那些副本才有资格竞选Leader 副本。假如出现这种情况:那些保存数据比较多的副本都挂了时,且此参数设置成 false,那么就坚持之前的原则,坚决不能让那些落后太多的副本竞选 Leader。这样做的后果是这个分区就不可用了,因为没有 Leader 了。反之如果是 true,那么 Kafka 允许你从那些“跑得慢”的副本中选一个出来当 Leader。这样做的后果是数据有可能就丢失了,一般默认值为false,为了防止万一可以显示指定为false。

  • auto.leader.rebalance.enable:是否允许定期进行 Leader 选举。

    此参数对生产环境影响非常大。设置它的值为 true 表示允许 Kafka 定期满足某些条件时地对一些 Topic 分区进行 Leader 重选举,换一次 Leader 代价很高的,原本向 A 发送请求的所有客户端都要切换成向 B 发送请求,而且这种换 Leader 本质上没有任何性能收益,因此建议在生产环境中把这个参数设置成 false。

  1. 关于数据留存方面的参数
  • log.retention.{hour|minutes|ms}

    这是个“三兄弟”,都是控制一条消息数据被保存多长时间。从优先级上来说 ms 设置最高、minutes 次之、hour 最低。通常情况下我们还是设置 hour 级别的多一些,比如log.retention.hour=168表示默认保存 7 天的数据,自动删除 7 天前的数据。

  • log.retention.bytes

    这是指定 Broker 为消息保存的总磁盘容量大小。这个值默认是 -1,表明你想在这台 Broker 上保存多少数据都可以,至少在容量方面 Broker 绝对为你开绿灯,不会做任何阻拦。这个参数真正发挥作用的场景其实是在云上构建多租户的 Kafka 集群:设想你要做一个云上的 Kafka 服务,每个租户只能使用 100GB 的磁盘空间,设置此参数可以避免个别“恶意”租户使用过多的磁盘空间。

  • message.max.bytes

    控制 Broker 能够接收的最大消息大小。默认的 1000012 太少了,还不到 1MB,在线上环境中设置一个比较大的值还是比较保险的做法。

Topic 级别参数

Kafka 也支持为不同的 Topic 设置不同的参数值。当前最新的 2.2 版本总共提供了大约 25 个 Topic 级别的参数,当然我们也不必全部了解它们的作用,这里我挑出了一些最关键的参数,你一定要把它们掌握清楚。

如果同时设置了 Topic 级别参数和全局 Broker 参数,到底听谁的呢?哪个说了算呢?

答案就是 Topic 级别参数会覆盖全局 Broker 参数的值,而每个 Topic 都能设置自己的参数值,这就是所谓的 Topic 级别参数。

举个例子说明一下,上一期我提到了消息数据的留存时间参数,在实际生产环境中,如果为所有 Topic 的数据都保存相当长的时间,这样做既不高效也无必要。更适当的做法是允许不同部门的 Topic 根据自身业务需要,设置自己的留存时间。如果只能设置全局 Broker 参数,那么势必要提取所有业务留存时间的最大值作为全局参数值,此时设置 Topic 级别参数把它覆盖,就是一个不错的选择。

下面我们依然按照用途分组的方式引出重要的 Topic 级别参数。

  1. 从保存消息方面来考量的话,下面这组参数是非常重要的:
  • retention.ms:规定了该 Topic 消息被保存的时长。默认是 7 天,即该 Topic 只保存最近 7 天的消息。一旦设置了这个值,它会覆盖掉 Broker 端的全局参数值。
  • retention.bytes:规定了要为该 Topic 预留多大的磁盘空间。和全局参数作用相似,这个值通常在多租户的 Kafka 集群中会有用武之地。当前默认值是 -1,表示可以无限使用磁盘空间。
  1. 从能处理的消息大小这个角度来看的话,有一个参数是必须要设置的,即
  • max.message.bytes

    它决定了 Kafka Broker 能够正常接收该 Topic 的最大消息大小。我知道目前在很多公司都把 Kafka 作为一个基础架构组件来运行,上面跑了很多的业务数据。如果在全局层面上,我们不好给出一个合适的最大消息值,那么不同业务部门能够自行设定这个 Topic 级别参数就显得非常必要了。在实际场景中,这种用法也确实是非常常见的。

怎么设置 Topic 级别参数

Topic 级别参数的设置就是这种情况,我们有两种方式可以设置:

  • 创建 Topic 时进行设置
  • 修改 Topic 时设置

我们先来看看如何在创建 Topic 时设置这些参数。我用上面提到的retention.msmax.message.bytes举例。设想你的部门需要将交易数据发送到 Kafka 进行处理,需要保存最近半年的交易数据,同时这些数据很大,通常都有几 MB,但一般不会超过 5MB。现在让我们用以下命令来创建 Topic:

bin/kafka-topics.sh--bootstrap-serverlocalhost:9092--create--topictransaction--partitions1--replication-factor1--configretention.ms=15552000000--configmax.message.bytes=5242880

我们只需要知道 Kafka 开放了kafka-topics命令供我们来创建 Topic 即可。对于上面这样一条命令,请注意结尾处的--config设置,我们就是在 config 后面指定了想要设置的 Topic 级别参数。

下面看看使用另一个自带的命令kafka-configs来修改 Topic 级别参数。假设我们现在要发送最大值是 10MB 的消息,该如何修改呢?命令如下:

 bin/kafka-configs.sh--zookeeperlocalhost:2181--entity-typetopics--entity-nametransaction--alter--add-configmax.message.bytes=10485760

总体来说,你只能使用这么两种方式来设置 Topic 级别参数。我个人的建议是,你最好始终坚持使用第二种方式来设置,并且在未来,Kafka 社区很有可能统一使用kafka-configs脚本来调整 Topic 级别参数。

JVM 参数

Kafka 服务器端代码是用 Scala 语言编写的,但终归还是编译成 Class 文件在 JVM 上运行,因此 JVM 参数设置对于 Kafka 集群的重要性不言而喻。

Kafka 自 2.0.0 版本开始,已经正式摒弃对 Java 7 的支持了,所以有条件的话至少使用 Java 8 吧。

说到 JVM 端设置,堆大小这个参数至关重要。虽然在后面我们还会讨论如何调优 Kafka 性能的问题,但现在我想无脑给出一个通用的建议:将你的 JVM 堆大小设置成 6GB 吧,这是目前业界比较公认的一个合理值。我见过很多人就是使用默认的 Heap Size 来跑 Kafka,说实话默认的 1GB 有点小,毕竟 Kafka Broker 在与客户端进行交互时会在 JVM 堆上创建大量的 ByteBuffer 实例,Heap Size 不能太小。

JVM 端配置的另一个重要参数就是垃圾回收器的设置,也就是平时常说的 GC 设置。如果你依然在使用 Java 7,那么可以根据以下法则选择合适的垃圾回收器:

  • 如果 Broker 所在机器的 CPU 资源非常充裕,建议使用 CMS 收集器。启用方法是指定

    -XX:+UseCurrentMarkSweepGC

  • 否则,使用吞吐量收集器。开启方法是指定

    -XX:+UseParallelGC

当然了,如果你已经在使用 Java 8 了,那么就用默认的 G1 收集器就好了。在没有任何调优的情况下,G1 表现得要比 CMS 出色,主要体现在更少的 Full GC,需要调整的参数更少等,所以使用 G1 就好了。

现在我们确定好了要设置的 JVM 参数,我们该如何为 Kafka 进行设置呢?有些奇怪的是,这个问题居然在 Kafka 官网没有被提及。其实设置的方法也很简单,你只需要设置下面这两个环境变量即可:

  • KAFKA_HEAP_OPTS:指定堆大小。
  • KAFKA_JVM_PERFORMANCE_OPTS:指定 GC 参数。

比如你可以这样启动 Kafka Broker,即在启动 Kafka Broker 之前,先设置上这两个环境变量:

$> export KAFKA_HEAP_OPTS=--Xms6g  --Xmx6g
$> export KAFKA_JVM_PERFORMANCE_OPTS= -server -XX:+UseG1GC -XX:MaxGCPauseMillis=20 -XX:InitiatingHeapOccupancyPercent=35 -XX:+ExplicitGCInvokesConcurrent -Djava.awt.headless=true
$> bin/kafka-server-start.sh config/server.properties

操作系统参数

最后我们来聊聊 Kafka 集群通常都需要设置哪些操作系统参数。通常情况下,Kafka 并不需要设置太多的 OS 参数,但有些因素最好还是关注一下,比如下面这几个:

  • 文件描述符限制
  • 文件系统类型
  • Swappiness
  • 提交时间

首先是ulimit -n。我觉得任何一个 Java 项目最好都调整下这个值。实际上,文件描述符系统资源并不像我们想象的那样昂贵,你不用太担心调大此值会有什么不利的影响。通常情况下将它设置成一个超大的值是合理的做法,比如ulimit -n 1000000。还记得电影《让子弹飞》里的对话吗:“你和钱,谁对我更重要?都不重要,没有你对我很重要!”。这个参数也有点这么个意思。其实设置这个参数一点都不重要,但不设置的话后果很严重,比如你会经常看到“Too many open files”的错误。

其次是文件系统类型的选择。这里所说的文件系统指的是如 ext3、ext4 或 XFS 这样的日志型文件系统。根据官网的测试报告,XFS 的性能要强于 ext4,所以生产环境最好还是使用 XFS。对了,最近有个 Kafka 使用 ZFS 的数据报告,貌似性能更加强劲,有条件的话不妨一试。

第三是 swap 的调优。网上很多文章都提到设置其为 0,将 swap 完全禁掉以防止 Kafka 进程使用 swap 空间。我个人反倒觉得还是不要设置成 0 比较好,我们可以设置成一个较小的值。为什么呢?因为一旦设置成 0,当物理内存耗尽时,操作系统会触发 OOM killer 这个组件,它会随机挑选一个进程然后 kill 掉,即根本不给用户任何的预警。但如果设置成一个比较小的值,当开始使用 swap 空间时,你至少能够观测到 Broker 性能开始出现急剧下降,从而给你进一步调优和诊断问题的时间。基于这个考虑,我个人建议将 swappniess 配置成一个接近 0 但不为 0 的值,比如 1。

最后是提交时间或者说是 Flush 落盘时间。向 Kafka 发送数据并不是真要等数据被写入磁盘才会认为成功,而是只要数据被写入到操作系统的页缓存(Page Cache)上就可以了,随后操作系统根据 LRU 算法会定期将页缓存上的“脏”数据落盘到物理磁盘上。这个定期就是由提交时间来确定的,默认是 5 秒。一般情况下我们会认为这个时间太频繁了,可以适当地增加提交间隔来降低物理磁盘的写操作。当然你可能会有这样的疑问:如果在页缓存中的数据在写入到磁盘前机器宕机了,那岂不是数据就丢失了。的确,这种情况数据确实就丢失了,但鉴于 Kafka 在软件层面已经提供了多副本的冗余机制,因此这里稍微拉大提交间隔去换取性能还是一个合理的做法。

小结

这里分享了关于 Kafka 集群设置的各类配置,包括 Topic 级别参数、JVM 参数以及操作系统参数。希望这些最佳实践能够在你搭建 Kafka 集群时助你一臂之力,但切记配置因环境而异,一定要结合自身业务需要以及具体的测试来验证它们的有效性。

Views: 298