You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何自动化从Elasticsearch Dev Tools提取数据的日常重复流程?

自动化Elasticsearch查询转CSV并导入Google Sheets的解决方案

Python脚本实现(快速上手,适合新手)

步骤1:安装依赖

pip install elasticsearch pandas gspread oauth2client

步骤2:编写完整脚本

from elasticsearch import Elasticsearch
import pandas as pd
import gspread
from oauth2client.service_account import ServiceAccountCredentials

# 1. 连接Elasticsearch(替换为你的ES地址、用户名密码)
es = Elasticsearch(
    ["https://your-es-host:9200"],
    basic_auth=("your-username", "your-password")
)

# 2. 执行你的查询DSL(替换为你实际的查询语句)
query_dsl = {
    "aggs": {
        "target_bucket": {
            # 这里替换为你实际的聚合逻辑,比如terms、date_histogram等
            "terms": {
                "field": "your-target-field.keyword"
            }
        }
    },
    "size": 0  # 仅返回聚合结果,无需原始文档
}

# 执行查询并获取响应
response = es.search(index="your-target-index", body=query_dsl)

# 3. 提取buckets数据
buckets = response["aggregations"]["target_bucket"]["buckets"]

# 4. 转换为结构化数据(无需第三方网站,直接用pandas处理)
df = pd.DataFrame(buckets)
# 可选:保存本地CSV文件
df.to_csv("es_buckets_result.csv", index=False)

# 5. 导入Google Sheets
# 提前准备:在Google Cloud Console创建服务账号,下载密钥JSON,将目标Sheets分享给服务账号邮箱
scope = ["https://spreadsheets.google.com/feeds", "https://www.googleapis.com/auth/drive"]
creds = ServiceAccountCredentials.from_json_keyfile_name("your-service-account-key.json", scope)
client = gspread.authorize(creds)

# 打开目标表格并写入数据
sheet = client.open("Your Looker Dashboard Sheet").sheet1
sheet.clear()  # 清空原有内容(可选)
sheet.update([df.columns.values.tolist()] + df.values.tolist())

关键提示

  • 替换所有your-xxx占位符为你的实际配置(ES地址、查询DSL、Sheets名称等)
  • Google Sheets API需要提前开启,服务账号必须拥有目标表格的编辑权限
  • 如果查询有多层嵌套聚合,只需调整buckets的提取路径即可

ETL工具调度(适合长期自动执行)

如果需要每天固定执行2-3次,推荐用Apache Airflow编排任务:

  1. 用Docker快速部署Airflow,安装所需Provider包:
pip install apache-airflow-providers-elasticsearch apache-airflow-providers-google
  1. 编写DAG调度文件:
from airflow import DAG
from airflow.providers.elasticsearch.hooks.elasticsearch import ElasticsearchHook
from airflow.providers.google.suite.hooks.sheets import GSheetsHook
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta
import pandas as pd

default_args = {
    'owner': 'luiz',
    'depends_on_past': False,
    'start_date': datetime(2024, 1, 1),
    'retries': 1,
    'retry_delay': timedelta(minutes=5)
}

def es_to_gsheets_etl():
    # 从Airflow连接池获取ES连接
    es_hook = ElasticsearchHook(elasticsearch_conn_id='your-es-connection')
    es = es_hook.get_conn()
    
    # 执行查询
    query_dsl = {
        # 替换为你的实际查询DSL
        "aggs": {"target_bucket": {"terms": {"field": "your-field.keyword"}}},
        "size": 0
    }
    response = es.search(index="your-index", body=query_dsl)
    buckets = response["aggregations"]["target_bucket"]["buckets"]
    
    # 转换数据
    df = pd.DataFrame(buckets)
    
    # 写入Google Sheets
    sheets_hook = GSheetsHook(gcp_conn_id='your-google-cloud-connection')
    sheets_hook.update_values(
        spreadsheet_id='your-sheets-id',
        range_name='Sheet1!A1',
        values=[df.columns.tolist()] + df.values.tolist()
    )

with DAG(
    'es_auto_to_gsheets',
    default_args=default_args,
    description='Automate ES query to Google Sheets for Looker Dashboard',
    schedule_interval='0 9,14,18 * * *',  # 每天9/14/18点执行,对应2-3次需求
    catchup=False,
) as dag:
    etl_task = PythonOperator(
        task_id='execute_etl',
        python_callable=es_to_gsheets_etl
    )

关键提示

  • 在Airflow UI的Admin > Connections中配置Elasticsearch和Google Cloud的连接信息
  • 调整schedule_interval可以自定义执行时间,支持cron表达式
  • Airflow自带日志监控、失败重试功能,适合长期稳定运行

内容的提问来源于stack exchange,提问作者Luiz Decio

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.06 12:05:42