Apache Druid接入Kafka实时流数据

简介: 一.任务配置文件使用类型为kafka{ "type": "kafka", "dataSchema": { "dimensionsSpec": {... ...}, "transformSpec":{.

一.任务配置文件

使用类型为kafka

{
  "type": "kafka",
  "dataSchema": {
      "dimensionsSpec": {... ...},
      "transformSpec":{... ...},
      "metricsSpec":{... ...}
  },
  "tuningConfig": {... ...},
  "ioConfig": {... ...}
}

Ⅰ).数据源

数据源配置部分5部分,表信息、解析器、数据转换、指标度量和聚合&查询粒度

  "dataSchema": {
    "dataSource": "druid_table_name",
    "parser": {},
    "transformSpec": {},
    "metricsSpec": {}
    }

a).表信息

"dataSource": "druid_table_name"

b).解析器

解析器包括:解析器类型、聚合字段和非聚合字段

    "parser": {
            "type": "string",
            "parseSpec": {
                "format": "json",
                "timestampSpec": {
                    "column": "time",
                    "format": "auto"
                },
                "dimensionsSpec": {
                    "dimensions": [
                        "appName",
                        "nodeName"
                    ],
                    "dimensionExclusions": []
                }
            }

c).数据转换

数据转换主要使用:时间转换、表达式、过滤等;其中,表达式可满足一些根据不同区间范围的指标统计类需求

    "transformSpec": {
            "transforms": [
                {
                    "type": "expression",
                    "name": "time",
                    "expression": "timestamp_format(time,'yyyy-MM-dd HH:mm:ss.SSS')"
                },
                {
                    "type": "expression",
                    "name": "status",
                    "expression": "if(status, 1, 0)"
                },
                {
                    "type": "expression",
                    "name": "processRange1",
                    "expression": "if(processTime<=100, 1, 0)"
                },
                {
                    "type": "expression",
                    "name": "processRange2",
                    "expression": "if(processTime>100  && processTime<= 500, 1, 0)"
                },
                {
                    "type": "expression",
                    "name": "processRange3",
                    "expression": "if(processTime>500, 1, 0)"
                }
            ]
        },

d).指标度量

指标度量主要使用:Sum、Max、Min、hyperUnique

    "metricsSpec": [
            {
                "name": "TotalTransCount",
                "fieldName": "count",
                "type": "longSum"
            },
            {
                "name": "MaxProcessTime",
                "fieldName": "processTime",
                "type": "longMax"
            },
            {
                "name": "MinProcessTime",
                "fieldName": "processTime",
                "type": "longMin"
            },
            {
                "name": "TotalProcessRange1",
                "fieldName": "processRange1",
                "type": "longSum"
            },
            {
                "name": "TotalProcessRange1",
                "fieldName": "processRange2",
                "type": "longSum"
            },
            {
                "name": "TotalProcessRange1",
                "fieldName": "processRange3",
                "type": "longSum"
            },
            {
                "name": "NodeName",
                "fieldName": "nodeName",
                "type": "hyperUnique"
            }
        ]

e).聚合&查询粒度

聚合&查询粒度主要使用:all、none、second、minute、fifteen_minute、thirty_minute、hour、day、week、month、quarter、year

      "granularitySpec": {
            "type": "uniform",
            "segmentGranularity": "DAY",
            "queryGranularity": "HOUR"
        }

Ⅱ).任务及Segments配置

任务及Segments主要配置生产segment的大小,合并任务进程数

    "tuningConfig": {
        "type": "kafka",
        "maxRowsPerSegment": 500000,
        "workerThreads": 2,
        "reportParseExceptions": false
    },

Ⅲ).数据接入信息

数据接入信息主要包括:kafka consumer相关配置和任务执行间隔

    "ioConfig": {
        "topic": "kafka_topic_name",
        "consumerProperties": {
            "bootstrap.servers": "hostname:9092"
        },
        "useEarliestOffset": false,
        "taskCount": 3,
        "replicas": 1,
        "taskDuration": "PT1H"
    }

Ⅳ).后聚合配置

主要使用在查询时,根据业务场景需求,需要配置在度量指标基础上运算获得二级指标

"aggregations":[
      {
         "type":"count",
         "name":"count"
      }
   ]

二.提交任务

任务配置文件kafka-streaming.json,提交命令如下

curl -XPOST -H'Content-Type: application/json' -d @quickstart/tutorial/kafka-streaming.json http://hostname:8081/druid/indexer/v1/supervisor

三.常遇问题

1.offsetOutofRangeException

kafka的offset过期清理 {“auto.offset.reset” : “latest” };导致 auto.offset.reset:None -- “None”表示offset过期不作任何处理,只抛出异常,即offsetOutofRangeException

解决办法:
    清空元信息库中druid_dataSource表, 表中记录了所有消费的Kafka topic对应的partition以及offset信息, 同时重启MiddleManager节

2.时差问题

默认:UTC,需要修改配置

a).各角色配置文件jvm.config

-Duser.timezone=UTC+0800

b).middleManager配置文件runtime.properties

# Task launch parameters
druid.indexer.runner.javaOpts=-server -Xmx8g -Duser.timezone=UTC+0800 -Dfile.encoding=UTF-8 -Djava.util.logging.manager=org.apache.logging.log4j.jul.LogManager

3.补跑数据

有时候任务失败,导致需要补录kafka数据,则修改数据接入中的配置信息

useEarliestOffset:true
目录
相关文章
|
26天前
|
消息中间件 前端开发 Kafka
【Azure 事件中心】使用Apache Flink 连接 Event Hubs 出错 Kafka error: No resolvable bootstrap urls
【Azure 事件中心】使用Apache Flink 连接 Event Hubs 出错 Kafka error: No resolvable bootstrap urls
|
15天前
|
存储 大数据 数据挖掘
【数据新纪元】Apache Doris:重塑实时分析性能,解锁大数据处理新速度,引爆数据价值潜能!
【9月更文挑战第5天】Apache Doris以其卓越的性能、灵活的架构和高效的数据处理能力,正在重塑实时分析的性能极限,解锁大数据处理的新速度,引爆数据价值的无限潜能。在未来的发展中,我们有理由相信Apache Doris将继续引领数据处理的潮流,为企业提供更快速、更准确、更智能的数据洞察和决策支持。让我们携手并进,共同探索数据新纪元的无限可能!
61 11
|
27天前
|
消息中间件 Java Kafka
【Azure 事件中心】在微软云中国区 (Mooncake) 上实验以Apache Kafka协议方式发送/接受Event Hubs消息 (Java版)
【Azure 事件中心】在微软云中国区 (Mooncake) 上实验以Apache Kafka协议方式发送/接受Event Hubs消息 (Java版)
|
22天前
|
消息中间件 Kafka 数据处理
实时数据流处理:Dask Streams 与 Apache Kafka 集成
【8月更文第29天】在现代数据处理领域,实时数据流处理已经成为不可或缺的一部分。随着物联网设备、社交媒体和其他实时数据源的普及,处理这些高吞吐量的数据流成为了一项挑战。Apache Kafka 作为一种高吞吐量的消息队列服务,被广泛应用于实时数据流处理场景中。Dask Streams 是 Dask 库的一个子模块,它为 Python 开发者提供了一个易于使用的实时数据流处理框架。本文将介绍如何将 Dask Streams 与 Apache Kafka 结合使用,以实现高效的数据流处理。
22 0
|
27天前
|
消息中间件 Java Kafka
【Azure 事件中心】开启 Apache Flink 制造者 Producer 示例代码中的日志输出 (连接 Azure Event Hub Kafka 终结点)
【Azure 事件中心】开启 Apache Flink 制造者 Producer 示例代码中的日志输出 (连接 Azure Event Hub Kafka 终结点)
|
27天前
|
消息中间件 Java Kafka
Kafka不重复消费的终极秘籍!解锁幂等性、偏移量、去重神器,让你的数据流稳如老狗,告别数据混乱时代!
【8月更文挑战第24天】Apache Kafka作为一款领先的分布式流处理平台,凭借其卓越的高吞吐量与低延迟特性,在大数据处理领域中占据重要地位。然而,在利用Kafka进行数据处理时,如何有效避免重复消费成为众多开发者关注的焦点。本文深入探讨了Kafka中可能出现重复消费的原因,并提出了四种实用的解决方案:利用消息偏移量手动控制消费进度;启用幂等性生产者确保消息不被重复发送;在消费者端实施去重机制;以及借助Kafka的事务支持实现精确的一次性处理。通过这些方法,开发者可根据不同的应用场景灵活选择最适合的策略,从而保障数据处理的准确性和一致性。
65 9
|
1月前
|
消息中间件 负载均衡 Java
"Kafka核心机制揭秘:深入探索Producer的高效数据发布策略与Java实战应用"
【8月更文挑战第10天】Apache Kafka作为顶级分布式流处理平台,其Producer组件是数据高效发布的引擎。Producer遵循高吞吐、低延迟等设计原则,采用分批发送、异步处理及数据压缩等技术提升性能。它支持按消息键值分区,确保数据有序并实现负载均衡;提供多种确认机制保证可靠性;具备失败重试功能确保消息最终送达。Java示例展示了基本配置与消息发送流程,体现了Producer的强大与灵活性。
52 3
|
22天前
|
消息中间件 存储 关系型数据库
实时计算 Flink版产品使用问题之如何使用Kafka Connector将数据写入到Kafka
实时计算Flink版作为一种强大的流处理和批处理统一的计算框架,广泛应用于各种需要实时数据处理和分析的场景。实时计算Flink版通常结合SQL接口、DataStream API、以及与上下游数据源和存储系统的丰富连接器,提供了一套全面的解决方案,以应对各种实时计算需求。其低延迟、高吞吐、容错性强的特点,使其成为众多企业和组织实时数据处理首选的技术平台。以下是实时计算Flink版的一些典型使用合集。
|
22天前
|
消息中间件 监控 Kafka
实时计算 Flink版产品使用问题之处理Kafka数据顺序时,怎么确保事件的顺序性
实时计算Flink版作为一种强大的流处理和批处理统一的计算框架,广泛应用于各种需要实时数据处理和分析的场景。实时计算Flink版通常结合SQL接口、DataStream API、以及与上下游数据源和存储系统的丰富连接器,提供了一套全面的解决方案,以应对各种实时计算需求。其低延迟、高吞吐、容错性强的特点,使其成为众多企业和组织实时数据处理首选的技术平台。以下是实时计算Flink版的一些典型使用合集。
|
26天前
|
消息中间件 缓存 Kafka
【Azure 事件中心】使用Kafka消费Azure EventHub中数据,遇见消费慢的情况可以如何来调节呢?
【Azure 事件中心】使用Kafka消费Azure EventHub中数据,遇见消费慢的情况可以如何来调节呢?

推荐镜像

更多