从 Apache Kafka 加载流数据
本教程展示了如何使用 Druid 的 Kafka 索引服务将数据从 Kafka 流加载到 Apache Druid 中。
本教程将指导您完成以下步骤:将来自 Koalas to the Max 游戏的嵌套点击流示例数据加载到 Kafka 主题中,然后将其摄取到 Druid 中。
先决条件
在按照本教程中的步骤操作之前,请按照快速入门中的说明下载 Druid,使用自动单机配置并在本地机器上运行。您无需预先加载任何数据。
下载并启动 Kafka
Apache Kafka 是一个高吞吐量的消息总线,与 Druid 配合使用效果很好。本教程请使用 Kafka 2.7.0。
-
要下载 Kafka,请在终端中运行以下命令
curl -O https://archive.apache.org/dist/kafka/2.7.0/kafka_2.13-2.7.0.tgz
tar -xzf kafka_2.13-2.7.0.tgz
cd kafka_2.13-2.7.0 -
如果您正在本教程使用的机器上运行 Kafka,请删除或重命名
/tmp目录下的kafka-logs文件夹。
Druid 和 Kafka 都依赖 Apache ZooKeeper 来协调和管理服务。由于 Druid 已经在运行,Kafka 启动时会自动连接到 Druid 的 ZooKeeper 实例。
在生产环境中,如果您在不同的机器上运行 Druid 和 Kafka,请在启动 Kafka 代理(broker)之前启动 Kafka ZooKeeper。
-
在 Kafka 根目录下,运行此命令以启动 Kafka 代理
./bin/kafka-server-start.sh config/server.properties -
在新的终端窗口中,导航到 Kafka 根目录并运行以下命令,创建一个名为
kttm的 Kafka 主题./bin/kafka-topics.sh --create --topic kttm --bootstrap-server localhost:9092成功添加主题后,Kafka 会返回一条消息:
Created topic kttm。
将数据加载到 Kafka 中
在本节中,您将把示例数据下载到教程目录,并将其发送到您的 Kafka 主题。
-
在您的 Kafka 根目录下,创建一个用于存放示例数据的目录
mkdir sample-data -
将示例数据下载到您的新目录并解压
(cd sample-data && curl -O https://static.imply.io/example-data/kttm-nested-v2/kttm-nested-v2-2019-08-25.json.gz) -
在您的 Kafka 根目录下,运行以下命令将示例事件发布到
kttmKafka 主题export KAFKA_OPTS="-Dfile.encoding=UTF-8"
gzcat ./sample-data/kttm-nested-v2-2019-08-25.json.gz | ./bin/kafka-console-producer.sh --broker-list localhost:9092 --topic kttm
将数据加载到 Druid 中
现在您的 Kafka 主题中已有数据,可以使用 Druid 的 Kafka 索引服务将数据摄取到 Druid 中。
为此,您可以使用 Druid 控制台的数据加载器(data loader),或者提交一个监督者规范(supervisor spec)。请按照以下步骤尝试每种方法。
使用控制台数据加载器加载数据
Druid 控制台数据加载器会引导您通过多个屏幕来配置监督者规范的每个部分,然后创建一个摄取任务来摄取 Kafka 数据。
要使用控制台数据加载器
-
导航至 localhost:8888 并点击 Load data(加载数据) > Streaming(流式)。

-
点击 Apache Kafka,然后点击 Connect data(连接数据)。
-
输入
localhost:9092作为引导服务器(bootstrap server),输入kttm作为主题,然后点击 Apply(应用),确保您能看到类似以下的数据
-
点击 Next: Parse data(下一步:解析数据)。

数据加载器会自动尝试确定数据的正确输入格式。对于示例数据,它会选择
json输入格式。您可以尝试不同的选项来预览 Druid 如何解析您的数据。 -
选择
json输入格式后,点击 Next: Parse time(下一步:解析时间)。您可能需要先点击 Apply(应用)。
Druid 的架构要求您指定一个主时间戳列。Druid 将时间戳存储在数据源的
__time列中。在生产环境中,如果数据中没有时间戳,您可以选择 Parse timestamp from:(解析时间戳来源:) 为None,从而使用占位符值。对于示例数据,数据加载器会自动将原始数据中的
timestamp列选为主要时间列。 -
点击三次 Next: ...(下一步:...),跳过 Transform(转换) 和 Filter(过滤) 步骤,到达 Configure schema(配置模式)。您无需在这两个步骤中输入任何内容,因为应用转换和过滤超出了本教程的范围。

-
在 Configure schema(配置模式) 步骤中,您可以为列选择数据类型,并配置要摄取到 Druid 中的维度和指标。控制台会自动完成大部分工作。请注意,维度
event、agent和geo_ip的类型均为json。 -
点击 Next: Partition(下一步:分区) 以配置 Druid 如何将数据分区到段(segments)中。

-
选择
day作为 Segment granularity(段粒度)。由于这是一个小型数据集,您不需要进行任何其他调整。点击 Next: Tune(下一步:调整) 以微调 Druid 摄取数据的方式。
-
在 Input tuning(输入调整) 中,将 Use earliest offset(使用最早偏移量) 设置为
True——这一点非常重要,因为您希望从流的开头开始消费数据。这里无需进行其他更改,因此点击 Next: Publish(下一步:发布)。
-
将数据源命名为
kttm-kafka,然后点击 Next: Edit spec(下一步:编辑规范) 以检查您的规范。
控制台会显示您构建的规范。您可以点击规范上方的按钮返回之前的步骤进行更改,并查看更改如何更新规范。您也可以直接编辑规范,并观察其在之前步骤中的变化。
-
点击 Submit(提交) 以创建摄取任务。
Druid 会显示任务视图,并将焦点锁定在新创建的监督者上。
任务视图会自动刷新,请等待监督者启动任务。当 Druid 开始摄取数据时,状态将从 Pending(待处理) 变为 Running(运行中)。

-
从标题导航至 Datasources(数据源) 视图。

当
kttm-kafka数据源出现在此处时,您就可以查询它了。详情请参阅 查询您的数据。
如果一分钟后数据源仍未出现,您可能没有设置监督者从流的开头读取数据——即 Tune(调整) 步骤中的 Use earliest offset 设置。请转到 Ingestion(摄取) 页面,使用 Actions(...)(操作) 菜单终止该监督者。重新加载示例数据,并在到达 Tune(调整) 步骤时应用正确的设置。
提交监督者规范
除了使用数据加载器,您还可以向 Druid 提交监督者规范。您可以在控制台中执行此操作,也可以使用 Druid API。
使用控制台
要使用 Druid 控制台提交监督者规范
-
点击控制台中的 Ingestion(摄取),然后点击刷新按钮旁边的省略号,选择 Submit JSON supervisor(提交 JSON 监督者)。
-
将此规范粘贴到 JSON 窗口中,然后点击 Submit(提交)。
{
"type": "kafka",
"spec": {
"ioConfig": {
"type": "kafka",
"consumerProperties": {
"bootstrap.servers": "localhost:9092"
},
"topic": "kttm",
"inputFormat": {
"type": "json"
},
"useEarliestOffset": true
},
"tuningConfig": {
"type": "kafka"
},
"dataSchema": {
"dataSource": "kttm-kafka-supervisor-console",
"timestampSpec": {
"column": "timestamp",
"format": "iso"
},
"dimensionsSpec": {
"dimensions": [
"session",
"number",
"client_ip",
"language",
"adblock_list",
"app_version",
"path",
"loaded_image",
"referrer",
"referrer_host",
"server_ip",
"screen",
"window",
{
"type": "long",
"name": "session_length"
},
"timezone",
"timezone_offset",
{
"type": "json",
"name": "event"
},
{
"type": "json",
"name": "agent"
},
{
"type": "json",
"name": "geo_ip"
}
]
},
"granularitySpec": {
"queryGranularity": "none",
"rollup": false,
"segmentGranularity": "day"
}
}
}
}这将启动监督者——监督者会生成开始监听传入数据的任务。
-
点击控制台主页上的 Tasks(任务) 以监控作业状态。此规范将
kttm主题中的数据写入名为kttm-kafka-supervisor-console的数据源中。
使用 API
您也可以使用 Druid API 提交监督者规范。
-
运行以下命令以下载示例规范
curl -o kttm-kafka-supervisor.json https://raw.githubusercontent.com/apache/druid/master/docs/assets/files/kttm-kafka-supervisor.json -
运行以下命令以提交
kttm-kafka-supervisor.json文件中的规范curl -X POST -H 'Content-Type: application/json' -d @kttm-kafka-supervisor.json https://:8081/druid/indexer/v1/supervisor在 Druid 成功创建监督者后,您将收到包含监督者 ID 的响应:
{"id":"kttm-kafka-supervisor-api"}。 -
点击控制台主页上的 Tasks(任务) 以监控作业状态。此规范将
kttm主题中的数据写入名为kttm-kafka-supervisor-api的数据源中。
查询您的数据
Druid 将数据发送到 Kafka 流后,该数据立即可供查询。点击 Druid 控制台中的 Query(查询),针对该数据源运行 SQL 查询。
由于本教程摄取的是小型数据集,您可以运行查询 SELECT * FROM "kttm-kafka" 来返回您创建的数据集中的所有数据。

查看查询数据教程,对新加载的数据运行一些示例查询。
延伸阅读
有关更多信息,请参阅以下主题
- Apache Kafka 摄取,了解有关从 Kafka 流加载数据以及维护 Druid Kafka 监督者的信息。