如何自动化从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编排任务:
- 用Docker快速部署Airflow,安装所需Provider包:
pip install apache-airflow-providers-elasticsearch apache-airflow-providers-google
- 编写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
相关产品推荐
相关产品推荐

