CSV离线导入:InfluxDB到KaiwuDB的兜底方案 文章目录每日一句正能量一、前言为什么需要CSV兜底二、环境准备2.1 源端环境InfluxDB2.2 目标端环境KaiwuDB2.3 中间服务器三、InfluxDB导出CSV3.1 使用influx命令导出3.2 使用Python导出3.3 导出数据验证四、CSV格式处理4.1 常见问题处理4.2 数据清洗脚本五、KaiwuDB导入CSV5.1 使用COPY命令导入5.2 使用Python导入5.3 批量导入优化六、性能优化6.1 导出优化6.2 导入优化6.3 性能对比七、常见问题7.1 导出问题7.2 导入问题八、总结每日一句正能量所有舒服的关系从不需要刻意强求全凭双向奔赴。与其耗尽心力强求一段不对等的关系不如识别并珍惜那些愿意与你同频奔赴的人。一、前言为什么需要CSV兜底前面几篇文章我分享了使用KDTS和DataX迁移InfluxDB到KaiwuDB的实战。有读者问“如果网络不通或者InfluxDB和KaiwuDB不在同一个网络怎么办”这是个好问题。在实际工作中我们确实遇到过这种情况网络隔离InfluxDB在内网KaiwuDB在外网中间隔着防火墙安全限制生产环境不允许部署第三方工具数据敏感设备监控数据不允许通过网络传输一次性迁移只需要迁移历史数据不需要实时同步这时候CSV离线导入就成了最后的保底方案。七个关键词离线迁移、CSV导出、格式转换、数据清洗、批量导入、性能优化、经验总结CSV方案的优势简单可靠不依赖任何工具纯文件操作可控性强可以人工审核数据确保敏感信息脱敏适用范围广任何能导出CSV的数据库都适用成本最低不需要额外安装软件当然CSV方案也有缺点速度较慢相比DataX需要手动处理编码、换行符等问题不适合超大规模数据千万级以上本文就把CSV离线导入的完整过程分享出来包括InfluxDB导出、CSV格式处理、KaiwuDB导入三个环节。二、环境准备2.1 源端环境InfluxDB# InfluxDB版本influx-version# InfluxDB shell version: 1.8.0# 数据库信息# 数据库名: sensor_db# Measurement: sensor_data# 数据量: 500GB# 数据点: 10亿2.2 目标端环境KaiwuDB# KaiwuDB版本kwbase version# KaiwuDB 3.2.0# 创建时序库CREATE DATABASE sensor_db;# 创建时序表CREATE TABLE sensor_data(timeTIMESTAMP NOT NULL, device_id STRING NOT NULL, location STRING, temperature FLOAT, humidity FLOAT, PRIMARY KEY(time, device_id));2.3 中间服务器用于数据中转需要足够的磁盘空间数据量的2-3倍能同时访问InfluxDB和KaiwuDB或分别访问安装有influx和kwbase客户端三、InfluxDB导出CSV3.1 使用influx命令导出InfluxDB提供了influx命令行工具可以导出数据为CSV格式。# 导出为CSVinflux-databasesensor_db-executeSELECT * FROM sensor_data-formatcsvsensor_data.csv# 导出指定时间范围influx-databasesensor_db-executeSELECT * FROM sensor_data WHERE time 2023-01-01 AND time 2024-01-01-formatcsvsensor_data_2023.csv3.2 使用Python导出如果influx命令行工具无法满足需求可以使用Python脚本导出。#!/usr/bin/env python3# export_influxdb.pyfrominfluxdbimportInfluxDBClientimportcsvdefexport_data():clientInfluxDBClient(hostlocalhost,port8086,databasesensor_db)# 查询数据resultclient.query(SELECT * FROM sensor_data)# 导出为CSVwithopen(sensor_data.csv,w,newline)asf:writercsv.writer(f)writer.writerow([time,device_id,location,temperature,humidity])forpointinresult.get_points():writer.writerow([point[time],point[device_id],point.get(location,),point[temperature],point[humidity]])if__name____main__:export_data()3.3 导出数据验证导出完成后需要验证数据的完整性和正确性# 查看文件大小ls-lhsensor_data.csv# 查看行数减去表头wc-lsensor_data.csv# 查看前10行head-10sensor_data.csv# 查看后10行tail-10sensor_data.csv四、CSV格式处理4.1 常见问题处理问题1字段中包含逗号如果字段内容中包含逗号会导致CSV格式错乱。解决方案使用双引号包裹字段。问题2字段中包含换行符如果字段内容中包含换行符如JSON字段会导致行数不对。解决方案在导出前替换换行符。#!/usr/bin/env python3# csv_clean.pyimportcsvdefclean_csv(input_file,output_file):withopen(input_file,r,newline)asf_in,open(output_file,w,newline)asf_out:readercsv.reader(f_in)writercsv.writer(f_out,quotingcsv.QUOTE_ALL)forrowinreader:# 处理每一行数据cleaned_row[]forfieldinrow:# 替换换行符cleaned_fieldfield.replace(\n,\\n)# 替换双引号cleaned_fieldcleaned_field.replace(,)cleaned_row.append(cleaned_field)writer.writerow(cleaned_row)if__name____main__:clean_csv(sensor_data_raw.csv,sensor_data_clean.csv)问题3编码问题InfluxDB默认使用UTF-8编码但有时候会出现编码问题。解决方案# 检查文件编码filesensor_data.csv# 转换编码iconv-futf-8-tutf-8 sensor_data.csvsensor_data_utf8.csv4.2 数据清洗脚本#!/usr/bin/env python3# csv_clean.pyimportcsvimportdatetimedefclean_csv(input_file,output_file):withopen(input_file,r,newline)asf_in,open(output_file,w,newline)asf_out:readercsv.reader(f_in)writercsv.writer(f_out,quotingcsv.QUOTE_ALL)# 写入表头headernext(reader)writer.writerow(header)forrowinreader:# 处理时间戳ifrow[0]:try:# 解析时间戳dtdatetime.datetime.fromisoformat(row[0].replace(Z,00:00))row[0]dt.strftime(%Y-%m-%d %H:%M:%S.%f)exceptValueError:pass# 处理空值foriinrange(len(row)):ifrow[i]orrow[i]isNone:row[i]NULLwriter.writerow(row)if__name____main__:clean_csv(sensor_data_raw.csv,sensor_data_clean.csv)五、KaiwuDB导入CSV5.1 使用COPY命令导入KaiwuDB支持COPY命令导入CSV文件。-- 连接到KaiwuDBkwbasesql--certs-dir/etc/kaiwudb/certs --host192.168.1.100:26257-- 使用COPY命令导入COPY sensor_data(time,device_id,location,temperature,humidity)FROM/tmp/sensor_data.csvWITH(FORMAT CSV,HEADERtrue,DELIMITER,);5.2 使用Python导入如果COPY命令不满足需求可以使用Python脚本导入。#!/usr/bin/env python3# import_kaiwudb.pyimportpsycopg2importcsvdefimport_data():connpsycopg2.connect(host192.168.1.100,port26257,databasesensor_db,userroot)cursorconn.cursor()withopen(sensor_data.csv,r)asf:readercsv.reader(f)next(reader)# 跳过表头forrowinreader:cursor.execute(INSERT INTO sensor_data (time, device_id, location, temperature, humidity) VALUES (%s, %s, %s, %s, %s),row)conn.commit()cursor.close()conn.close()if__name____main__:import_data()5.3 批量导入优化对于大数据量建议使用批量导入。#!/usr/bin/env python3# batch_import.pyimportpsycopg2importcsvdefbatch_import(batch_size1000):connpsycopg2.connect(host192.168.1.100,port26257,databasesensor_db,userroot)cursorconn.cursor()withopen(sensor_data.csv,r)asf:readercsv.reader(f)next(reader)# 跳过表头batch[]forrowinreader:batch.append(row)iflen(batch)batch_size:cursor.executemany(INSERT INTO sensor_data (time, device_id, location, temperature, humidity) VALUES (%s, %s, %s, %s, %s),batch)conn.commit()batch[]# 处理最后一批ifbatch:cursor.executemany(INSERT INTO sensor_data (time, device_id, location, temperature, humidity) VALUES (%s, %s, %s, %s, %s),batch)conn.commit()cursor.close()conn.close()if__name____main__:batch_import()六、性能优化6.1 导出优化# 使用并行导出influx-databasesensor_db-executeSELECT * FROM sensor_data WHERE time 2023-01-01 AND time 2023-02-01-formatcsvsensor_data_2023_01.csvinflux-databasesensor_db-executeSELECT * FROM sensor_data WHERE time 2023-02-01 AND time 2023-03-01-formatcsvsensor_data_2023_02.csvwait# 使用压缩gzipsensor_data.csv6.2 导入优化-- 禁用索引ALTERTABLEsensor_dataDISABLEINDEX;-- 导入数据COPY sensor_dataFROM/tmp/sensor_data.csvWITH(FORMAT CSV,HEADERtrue);-- 启用索引ALTERTABLEsensor_dataENABLEINDEX;-- 分析表ANALYZEsensor_data;6.3 性能对比方案500万数据耗时优点缺点COPY命令15分钟简单快速需要处理空值单条INSERT5小时灵活可控速度太慢批量INSERT30分钟速度较快需要编写脚本LOAD DATA10分钟速度最快需要文件权限七、常见问题7.1 导出问题问题1导出数据量太大内存不足解决方案# 分批导出influx-databasesensor_db-executeSELECT * FROM sensor_data WHERE time 2023-01-01 AND time 2023-02-01-formatcsvsensor_data_2023_01.csv influx-databasesensor_db-executeSELECT * FROM sensor_data WHERE time 2023-02-01 AND time 2023-03-01-formatcsvsensor_data_2023_02.csv问题2导出数据格式不对解决方案#!/usr/bin/env python3# export_format.pyfrominfluxdbimportInfluxDBClientimportcsvdefexport_data():clientInfluxDBClient(hostlocalhost,port8086,databasesensor_db)# 查询数据resultclient.query(SELECT * FROM sensor_data)# 导出为CSVwithopen(sensor_data.csv,w,newline)asf:writercsv.writer(f,quotingcsv.QUOTE_ALL)writer.writerow([time,device_id,location,temperature,humidity])forpointinresult.get_points():writer.writerow([point[time],point[device_id],point.get(location,),point[temperature],point[humidity]])if__name____main__:export_data()7.2 导入问题问题3导入时报错invalid input syntax原因CSV文件中的空值被解析为空字符串而KaiwuDB期望NULL。解决方案-- 使用NULL参数指定空值COPY sensor_data(time,device_id,location,temperature,humidity)FROM/tmp/sensor_data.csvWITH(FORMAT CSV,HEADERtrue,DELIMITER,,NULL);问题4导入速度慢解决方案增大批次大小5000-10000禁用索引使用批量导入八、总结CSV离线导入虽然速度不如DataX但在特定场景下是最可靠的方案。CSV方案的优点不依赖网络适合隔离环境不依赖第三方工具安全性高可以人工审核数据确保数据质量成本低不需要额外软件CSV方案的缺点速度较慢不适合大规模数据需要手动处理编码、换行符等问题需要编写脚本处理复杂场景最佳实践导出前确认编码、字段分隔符、换行符导出后验证数据完整性导入前清洗数据处理空值、特殊字符导入后校验数据确保一致性保留备份导入完成后保留CSV文件便于回滚如果你正在考虑InfluxDB迁移到KaiwuDB建议先评估网络环境和数据规模再选择合适的迁移方案。转载自https://blog.csdn.net/u014727709/article/details/164097089欢迎 点赞✍评论⭐收藏欢迎指正