本文目录导读:

- 使用
curl+ Shell 脚本(最适合:少量 JSON/CSV 文件,一次性导入) - 使用
Logstash(最主流、最稳定、企业级) - 使用 Python 脚本(最灵活,适合复杂逻辑)
- 使用
Filebeat(最适合:实时日志文件) - 总结:如何选择?
- 一个关键的“实用”提醒
是的,有很多实用脚本可以自动导入数据到 Elasticsearch (ES),具体采用哪种脚本,取决于你的数据源(如 CSV、JSON、数据库、日志文件)和导入频率(一次性、定时增量)。
以下介绍 4 种最主流、最实用的脚本方案,从简单到强大依次排列:
使用 curl + Shell 脚本(最适合:少量 JSON/CSV 文件,一次性导入)
这是最轻量级的方式,无需安装任何额外客户端。
场景: 你有一个 data.json 文件,希望导入 ES。
#!/bin/bash
# bulk_import.sh
# 使用 ES Bulk API 批量导入
ES_HOST="http://localhost:9200"
INDEX_NAME="my_index"
FILE_PATH="./data.json"
# 注意:Bulk API 要求每行数据格式为:action 行 + document 行
# 你可以用 jq 对普通 JSON 进行转换,或者预先准备好多行格式
echo "开始导入数据到 ES..."
# 方法1:直接 Bulk 导入(如果文件已经是 Bulk 格式)
# curl -X POST "$ES_HOST/$INDEX_NAME/_bulk" \
# -H "Content-Type: application/json" \
# --data-binary @"$FILE_PATH"
# 方法2:从普通 JSON 数组转换并导入(需要 jq 工具)
# 假设你的 data.json 是 [{"id":1,"name":"A"}, {"id":2,"name":"B"}]
if command -v jq &> /dev/null; then
jq -c '.[] | {"index": {"_index": "'$INDEX_NAME'"}}, .' "$FILE_PATH" | \
curl -X POST "$ES_HOST/_bulk" \
-H "Content-Type: application/json" \
--data-binary @-
echo "导入完成!"
else
echo "错误:需要安装 jq (brew install jq / apt install jq)"
fi
运行: chmod +x bulk_import.sh && ./bulk_import.sh
使用 Logstash(最主流、最稳定、企业级)
Logstash 是 ELK 栈的原生数据采集工具,支持多种输入(文件、数据库、Kafka、HTTP)和强大过滤。
场景: 每天凌晨从 MySQL 数据库增量导入数据到 ES。
配置文件 mysql-to-es.conf:
input {
jdbc {
jdbc_driver_library => "/path/to/mysql-connector-java-8.0.22.jar"
jdbc_driver_class => "com.mysql.cj.jdbc.Driver"
jdbc_connection_string => "jdbc:mysql://localhost:3306/mydb"
jdbc_user => "root"
jdbc_password => "password"
# 核心:使用 tracking_column 实现增量同步
statement => "SELECT * FROM orders WHERE update_time > :sql_last_value"
use_column_value => true
tracking_column => "update_time"
tracking_column_type => "timestamp"
# 第一次跑会拉取全部,以后只拉取更新的
schedule => "cron(0 2 * * *)" # 每天凌晨2点执行
}
}
filter {
# 可选:数据清洗、类型转换
mutate {
convert => { "price" => "float" }
rename => { "id" => "order_id" }
}
}
output {
elasticsearch {
hosts => ["localhost:9200"]
index => "orders-%{+YYYY.MM.dd}" # 按日期建索引
document_id => "%{order_id}" # 防止重复
}
stdout { codec => rubydebug } # 打印调试信息
}
运行命令:
logstash -f mysql-to-es.conf
优点: 稳定、有重试机制、数据不丢失、可处理各种格式。 缺点: 需要安装 Java,配置相对复杂。
使用 Python 脚本(最灵活,适合复杂逻辑)
Python 的 elasticsearch-py 库非常强大,你可以根据自己的业务逻辑编写脚本。
场景: 读取 CSV 文件,进行字段映射,然后写入 ES。
# import_to_es.py
import pandas as pd
from elasticsearch import Elasticsearch, helpers
import chardet # 处理中文编码
# 连接 ES
es = Elasticsearch(['http://localhost:9200'])
# 读取 CSV
def detect_encoding(file_path):
with open(file_path, 'rb') as f:
result = chardet.detect(f.read(10000))
return result['encoding']
df = pd.read_csv('data.csv', encoding=detect_encoding('data.csv'))
df = df.fillna('') # 处理空值
# 生成 Bulk 操作
def generate_actions(df, index_name):
for idx, row in df.iterrows():
doc = row.to_dict()
yield {
"_index": index_name,
"_id": doc.get('id', idx), # 用 id 字段或行号作为文档 ID
"_source": doc
}
# 执行批量导入
success, errors = helpers.bulk(
es,
generate_actions(df, "my_csv_index"),
chunk_size=500, # 每批 500 条
request_timeout=30
)
print(f"成功导入 {success} 条数据,失败 {len(errors) if errors else 0} 条")
运行: pip install elasticsearch pandas chardet && python import_to_es.py
扩展: 如果你要从 API 抓数据,只需把 df 替换为 requests.get(url).json() 即可。
使用 Filebeat(最适合:实时日志文件)
如果你的数据是不断追加的日志文件(如 access.log),Filebeat 是首选。
场景: 监控 /var/log/nginx/access.log,实时传给 ES。
配置文件 filebeat.yml:
filebeat.inputs:
- type: log
enabled: true
paths:
- /var/log/nginx/access.log
output.elasticsearch:
hosts: ["http://localhost:9200"]
index: "nginx-access-%{+yyyy.MM.dd}"
setup.kibana:
host: "http://localhost:5601"
运行:
filebeat -e -c filebeat.yml
优点: 极轻量、内存占用低、即插即用。 缺点: 主要用于行文本日志,不适合复杂结构化数据。
如何选择?
| 你的场景 | 推荐方案 | 原因 |
|---|---|---|
| 一次性的小文件 ( < 100MB) | Shell + curl | 最简单,无需装任何东西 |
| CSV / JSON 定期导入 | Python 脚本 | 灵活可控,可做数据清洗 |
| MySQL / Oracle 增量同步 | Logstash (JDBC) | 专门干这个的,稳定可靠 |
| 服务器日志实时采集 | Filebeat | 轻量,高性能,日志首选 |
| 需要复杂ETL处理 | Logstash 或 Python | 支持 grok、filter 等高级转换 |
一个关键的“实用”提醒
无论你用哪种脚本,对于大规模导入,请务必记住这 3 点才能保证“实用”(否则会挂):
- 使用
_bulkAPI:不要一条一条 insert,批量导入(500-1000条一批)速度提升 100 倍。 - 关闭 Refresh 和 Replicas:导入前设置:
PUT /my_index/_settings { "index": { "refresh_interval": "-1", # 导入完再开启 "number_of_replicas": 0 # 导入完再增加副本 } } - 导入完再恢复:
PUT /my_index/_settings { "index": { "refresh_interval": "1s", "number_of_replicas": 1 } }
如果你能告诉我你的具体数据来源(每天从数据库导出 CSV”或“实时爬虫数据”),我可以给你更精准的脚本。