python腳本將mysql數(shù)據(jù)寫入doris過(guò)程
python腳本將mysql數(shù)據(jù)寫入doris
flink-cdc連接的FE,數(shù)據(jù)寫入時(shí)正常的
但通過(guò)python腳本使用stream load方式直接連接FE,提示
NOT_AUTHORIZED]no valid Basic authorization
如果改用直接連接BE節(jié)點(diǎn),寫入又提示
[DATA_QUALITY_ERROR]too many filtered rows
這個(gè)就是搞技術(shù)最麻煩的地方,提示的異常,有沒有具體信息,根本就不知道原因是啥,要猜。
from sqlalchemy import create_engine
import pandas as pd
import math
import requests
from loguru import logger
import gzip
import io
import uuid
from requests.auth import HTTPBasicAuth
import json
# MySQL 配置
ziping_config = {
'user': 'root',
'password': '123456',
'host': 'localhost',
'port': 3306,
'database': 'ziping'
}
# 創(chuàng)建 SQLAlchemy 引擎
# ziping_engine = create_engine(
# "mysql+pymysql://{user}:{password}@{host}:{port}/{database}?charset=utf8mb4".format(**ziping_config)
# )
def sync_zp_bazi_info():
# 從 MySQL 讀取數(shù)據(jù)
ziping_engine = create_engine(
"mysql+pymysql://{user}:{password}@{host}:{port}/{database}".format(**ziping_config)
,connect_args={'charset': 'utf8mb4'}
)
logger.info("連接 MySQL 成功")
query = "SELECT id, year, month, day, hour, year_xun, month_xun, day_xun, hour_xun, " \
"yg, yz, mg, mz, dg, dz, hg, hz, year_shu, month_shu, day_shu, hour_shu, yz_cs, " \
"mz_cs, dz_cs, hz_cs, yg_sk, yz_sk, mg_sk, mz_sk, dz_sk, hg_sk, hz_sk, year_ny, " \
"month_ny, day_ny, hour_ny, jin_cn, mu_cn, shui_cn, huo_cn, tu_cn, zy_cn, py_cn, " \
"zc_cn, pc_cn, zg_cn, qs_cn, ss_cn, sg_cn, bj_cn, jc_cn, ym_gj, md_gj, dh_gj FROM zp_bazi_info"
df = pd.read_sql(query, ziping_engine)
logger.info("讀取數(shù)據(jù)成功")
# Stream Load 配置
stream_load_url = "http://10.101.1.31:8040/api/fay/zp_bazi_info/_stream_load"
auth = HTTPBasicAuth('root', '123456')
headers = {
"Content-Encoding":"gzip",
"Content-Type": "text/csv; charset=UTF-8", # 數(shù)據(jù)格式為 CSV
"Expect": "100-continue", # 支持大文件傳輸
}
# 分批次大小
batch_size = 100000 # 每批次 10 萬(wàn)條
total_rows = len(df)
num_batches = math.ceil(total_rows / batch_size)
# 分批次導(dǎo)入
for i in range(num_batches):
start = i * batch_size
end = min((i + 1) * batch_size, total_rows)
batch_df = df[start:end]
# 將批次數(shù)據(jù)轉(zhuǎn)換為 CSV
batch_csv = batch_df.to_csv(index=False, header=False,sep='\t').encode('utf-8')
headers["Content-Length"] = str(len(batch_csv))
# 壓縮 CSV
# compressed_data = compress_csv(batch_csv)
# 發(fā)送 Stream Load 請(qǐng)求
headers["label"] = str(uuid.uuid1())
response = requests.put(stream_load_url, headers=headers, data=batch_csv,auth=auth,timeout=60)
# 檢查導(dǎo)入結(jié)果
result = response.json()
if result['Status'] != 'Fail':
logger.info(result['Message'])
logger.info(f"批次 {i + 1}/{num_batches} 導(dǎo)入成功")
else:
logger.info(f"批次 {i + 1}/{num_batches} 導(dǎo)入失敗")
logger.info("錯(cuò)誤信息:{}", result['Message'])
break
def compress_csv(batch_csv):
buffer = io.BytesIO()
with gzip.GzipFile(fileobj=buffer, mode='wb') as f:
f.write(batch_csv)
compressed_data = buffer.getvalue()
return compressed_data
if __name__ == '__main__':
sync_zp_bazi_info()
從BE中可以查看到日志,csv文件中到doris中被解析為1列了。
看來(lái)問(wèn)題出在batch_df.to_csv(index=False, header=False,sep='\t'),需要增加sep='\t'

這是將mysql的表結(jié)構(gòu)直接轉(zhuǎn)成doris的表
CREATE TABLE `zp_bazi_info` (
`id` CHAR(11) NOT NULL COMMENT "ID",
`year` CHAR(2) NOT NULL COMMENT "年柱",
`month` CHAR(2) COMMENT "月柱",
`day` CHAR(2) COMMENT "日柱",
`hour` CHAR(2) COMMENT "時(shí)柱",
`year_xun` CHAR(2) COMMENT "年柱旬",
`month_xun` CHAR(2) COMMENT "月柱旬",
`day_xun` CHAR(2) COMMENT "日柱旬",
`hour_xun` CHAR(2) COMMENT "時(shí)柱旬",
`yg` CHAR(1) COMMENT "年干",
`yz` CHAR(1) COMMENT "年支",
`mg` CHAR(1) COMMENT "月干",
`mz` CHAR(1) COMMENT "月支",
`dg` CHAR(1) COMMENT "日干",
`dz` CHAR(1) COMMENT "日支",
`hg` CHAR(1) COMMENT "時(shí)干",
`hz` CHAR(1) COMMENT "時(shí)支",
`year_shu` TINYINT COMMENT "年柱數(shù)",
`month_shu` TINYINT COMMENT "月柱數(shù)",
`day_shu` TINYINT COMMENT "日柱數(shù)",
`hour_shu` TINYINT COMMENT "時(shí)柱數(shù)",
`yz_cs` CHAR(2) COMMENT "年長(zhǎng)生",
`mz_cs` CHAR(2) COMMENT "月長(zhǎng)生",
`dz_cs` CHAR(2) COMMENT "日長(zhǎng)生",
`hz_cs` CHAR(2) COMMENT "時(shí)長(zhǎng)生",
`yg_sk` CHAR(1) COMMENT "年干生克",
`yz_sk` CHAR(1) COMMENT "年支生克",
`mg_sk` CHAR(1) COMMENT "月干生克",
`mz_sk` CHAR(1) COMMENT "月支生克",
`dz_sk` CHAR(1) COMMENT "日支生克",
`hg_sk` CHAR(1) COMMENT "時(shí)干生克",
`hz_sk` CHAR(1) COMMENT "時(shí)支生克",
`year_ny` CHAR(3) COMMENT "年納音",
`month_ny` CHAR(3) COMMENT "月納音",
`day_ny` CHAR(3) COMMENT "日納音",
`hour_ny` CHAR(3) COMMENT "時(shí)納音",
`jin_cn` TINYINT COMMENT "金數(shù)量",
`mu_cn` TINYINT COMMENT "木數(shù)量",
`shui_cn` TINYINT COMMENT "水?dāng)?shù)量",
`huo_cn` TINYINT COMMENT "火數(shù)量",
`tu_cn` TINYINT COMMENT "土數(shù)量",
`zy_cn` TINYINT COMMENT "正印數(shù)量",
`py_cn` TINYINT COMMENT "偏印數(shù)量",
`zc_cn` TINYINT COMMENT "正財(cái)數(shù)量",
`pc_cn` TINYINT COMMENT "偏財(cái)數(shù)量",
`zg_cn` TINYINT COMMENT "正官數(shù)量",
`qs_cn` TINYINT COMMENT "七殺數(shù)量",
`ss_cn` TINYINT COMMENT "食神數(shù)量",
`sg_cn` TINYINT COMMENT "傷官數(shù)量",
`bj_cn` TINYINT COMMENT "比肩數(shù)量",
`jc_cn` TINYINT COMMENT "劫財(cái)數(shù)量",
`ym_gj` CHAR(1) COMMENT "年月拱夾",
`md_gj` CHAR(1) COMMENT "月時(shí)拱夾",
`dh_gj` CHAR(1) COMMENT "日時(shí)拱夾"
)
ENGINE=OLAP
DUPLICATE KEY(`id`)
DISTRIBUTED BY HASH(`id`) BUCKETS 10
PROPERTIES (
"replication_num" = "1"
);
但實(shí)際寫入,報(bào)下面的錯(cuò)誤,看來(lái)錯(cuò)誤原因應(yīng)該是mysql的字節(jié)與doris的字節(jié)計(jì)算不一樣。
Reason: column_name[id], the length of input is too long than schema. first 32 bytes of input str: [丁丑 丁未 丁丑 丁未] schema length: 11; actual length: 27; . src line [];
?因?yàn)閕d為主鍵,因此執(zhí)行下面的語(yǔ)句,會(huì)提示
Execution failed: Error Failed to execute sql: java.sql.SQLException: (conn=24) errCode = 2, detailMessage = No key column left. index[zp_bazi_info]
ALTER TABLE zp_bazi_info MODIFY COLUMN `id` VARCHAR(32) NOT NULL COMMENT "ID";
這里看一下,flink-cdc是怎么做的。

你可以看到varchar在doris中變成了4倍
utf8mb4 編碼(這是 MySQL 默認(rèn)推薦的 UTF-8 實(shí)現(xiàn),支持完整的 Unicode 字符集,包括表情符號(hào)等),flink采取了最壞的保守策略。而char是保持不變。

CREATE TABLE `zp_bazi_info` (
`id` VARCHAR(27) NOT NULL COMMENT "ID",
`year` VARCHAR(6) NOT NULL COMMENT "年柱",
`month` VARCHAR(6) COMMENT "月柱",
`day` VARCHAR(6) COMMENT "日柱",
`hour` VARCHAR(6) COMMENT "時(shí)柱",
`year_xun` VARCHAR(6) COMMENT "年柱旬",
`month_xun` VARCHAR(6) COMMENT "月柱旬",
`day_xun` VARCHAR(6) COMMENT "日柱旬",
`hour_xun` VARCHAR(6) COMMENT "時(shí)柱旬",
`yg` VARCHAR(9) COMMENT "年干",
`yz` VARCHAR(9) COMMENT "年支",
`mg` VARCHAR(9) COMMENT "月干",
`mz` VARCHAR(9) COMMENT "月支",
`dg` VARCHAR(9) COMMENT "日干",
`dz` VARCHAR(9) COMMENT "日支",
`hg` VARCHAR(9) COMMENT "時(shí)干",
`hz` VARCHAR(9) COMMENT "時(shí)支",
`year_shu` TINYINT COMMENT "年柱數(shù)",
`month_shu` TINYINT COMMENT "月柱數(shù)",
`day_shu` TINYINT COMMENT "日柱數(shù)",
`hour_shu` TINYINT COMMENT "時(shí)柱數(shù)",
`yz_cs` VARCHAR(6) COMMENT "年長(zhǎng)生",
`mz_cs` VARCHAR(6) COMMENT "月長(zhǎng)生",
`dz_cs` VARCHAR(6) COMMENT "日長(zhǎng)生",
`hz_cs` VARCHAR(6) COMMENT "時(shí)長(zhǎng)生",
`yg_sk` VARCHAR(9) COMMENT "年干生克",
`yz_sk` VARCHAR(9) COMMENT "年支生克",
`mg_sk` VARCHAR(9) COMMENT "月干生克",
`mz_sk` VARCHAR(9) COMMENT "月支生克",
`dz_sk` VARCHAR(9) COMMENT "日支生克",
`hg_sk` VARCHAR(9) COMMENT "時(shí)干生克",
`hz_sk` VARCHAR(9) COMMENT "時(shí)支生克",
`year_ny` VARCHAR(9) COMMENT "年納音",
`month_ny` VARCHAR(9) COMMENT "月納音",
`day_ny` VARCHAR(9) COMMENT "日納音",
`hour_ny` VARCHAR(9) COMMENT "時(shí)納音",
`jin_cn` TINYINT COMMENT "金數(shù)量",
`mu_cn` TINYINT COMMENT "木數(shù)量",
`shui_cn` TINYINT COMMENT "水?dāng)?shù)量",
`huo_cn` TINYINT COMMENT "火數(shù)量",
`tu_cn` TINYINT COMMENT "土數(shù)量",
`zy_cn` TINYINT COMMENT "正印數(shù)量",
`py_cn` TINYINT COMMENT "偏印數(shù)量",
`zc_cn` TINYINT COMMENT "正財(cái)數(shù)量",
`pc_cn` TINYINT COMMENT "偏財(cái)數(shù)量",
`zg_cn` TINYINT COMMENT "正官數(shù)量",
`qs_cn` TINYINT COMMENT "七殺數(shù)量",
`ss_cn` TINYINT COMMENT "食神數(shù)量",
`sg_cn` TINYINT COMMENT "傷官數(shù)量",
`bj_cn` TINYINT COMMENT "比肩數(shù)量",
`jc_cn` TINYINT COMMENT "劫財(cái)數(shù)量",
`ym_gj` VARCHAR(9) COMMENT "年月拱夾",
`md_gj` VARCHAR(9) COMMENT "月時(shí)拱夾",
`dh_gj` VARCHAR(9) COMMENT "日時(shí)拱夾"
)
ENGINE=OLAP
UNIQUE KEY(`id`)
DISTRIBUTED BY HASH(`id`) BUCKETS 10
PROPERTIES (
"replication_num" = "1"
);
擴(kuò)展后,數(shù)據(jù)寫入正常了,于是我又驗(yàn)證了,反復(fù)執(zhí)行,看看有沒有問(wèn)題。結(jié)果在doris中出現(xiàn)了兩條數(shù)據(jù)。id不是key,為什么會(huì)重復(fù)寫入呢?
因?yàn)樵贒oris中,“Duplicate Key”是一個(gè)冗余模型的特性。這個(gè)模型的數(shù)據(jù)完全按照導(dǎo)入文件中的數(shù)據(jù)進(jìn)行存儲(chǔ),不會(huì)有任何聚合。即使兩行數(shù)據(jù)完全相同,也都會(huì)保留。
這個(gè)Duplicate模型針對(duì)日志是可以
但是針對(duì)我們的系統(tǒng)表是不合適的,這里就需要采用UNIQUE KEY

51萬(wàn)數(shù)據(jù),寫入過(guò)程日志,doris的數(shù)據(jù)寫入速度還挺快。
2025-03-01 15:27:18.756 | INFO | __main__:sync_zp_bazi_info:33 - 連接 MySQL 成功 2025-03-01 15:29:06.586 | INFO | __main__:sync_zp_bazi_info:40 - 讀取數(shù)據(jù)成功 2025-03-01 15:29:21.388 | INFO | __main__:sync_zp_bazi_info:74 - OK 2025-03-01 15:29:21.388 | INFO | __main__:sync_zp_bazi_info:75 - 批次 1/6 導(dǎo)入成功 2025-03-01 15:29:36.447 | INFO | __main__:sync_zp_bazi_info:74 - OK 2025-03-01 15:29:36.447 | INFO | __main__:sync_zp_bazi_info:75 - 批次 2/6 導(dǎo)入成功 2025-03-01 15:29:52.508 | INFO | __main__:sync_zp_bazi_info:74 - OK 2025-03-01 15:29:52.509 | INFO | __main__:sync_zp_bazi_info:75 - 批次 3/6 導(dǎo)入成功 2025-03-01 15:30:06.747 | INFO | __main__:sync_zp_bazi_info:74 - OK 2025-03-01 15:30:06.747 | INFO | __main__:sync_zp_bazi_info:75 - 批次 4/6 導(dǎo)入成功 2025-03-01 15:30:22.621 | INFO | __main__:sync_zp_bazi_info:74 - OK 2025-03-01 15:30:22.621 | INFO | __main__:sync_zp_bazi_info:75 - 批次 5/6 導(dǎo)入成功 2025-03-01 15:30:24.921 | INFO | __main__:sync_zp_bazi_info:74 - OK 2025-03-01 15:30:24.921 | INFO | __main__:sync_zp_bazi_info:75 - 批次 6/6 導(dǎo)入成功 Process finished with exit code 0
執(zhí)行
select * from zp_bazi_info where year='乙丑' and month='戊子'
mysql需要2.232s,而doris卻只需要92ms,doris這個(gè)查詢比mysql快了24倍。那么為什么doris那么快呢?
總結(jié)
以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。
相關(guān)文章
pandas中的DataFrame按指定順序輸出所有列的方法
下面小編就為大家分享一篇pandas中的DataFrame按指定順序輸出所有列的方法,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧2018-04-04
Django的ListView超詳細(xì)用法(含分頁(yè)paginate)
這篇文章主要介紹了Django的ListView超詳細(xì)用法(含分頁(yè)paginate),文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧2020-05-05
Django2.1集成xadmin管理后臺(tái)所遇到的錯(cuò)誤集錦(填坑)
這篇文章主要介紹了Django2.1集成xadmin管理后臺(tái)所遇到的錯(cuò)誤集錦(填坑),小編覺得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧2018-12-12
python開發(fā)實(shí)例之Python的Twisted框架中Deferred對(duì)象的詳細(xì)用法與實(shí)例
這篇文章主要介紹了python開發(fā)實(shí)例之Python的Twisted框架中Deferred對(duì)象的詳細(xì)用法與實(shí)例,需要的朋友可以參考下2020-03-03
Python自動(dòng)化處理Word文檔表格格式的方法技巧
本項(xiàng)目介紹了如何使用Python及其 python-docx 庫(kù)來(lái)設(shè)置和修改Microsoft Word文檔中的表格格式,這對(duì)于需要自動(dòng)化報(bào)告生成、數(shù)據(jù)分析或批量定制文檔模板的場(chǎng)景尤為有用,通過(guò)實(shí)際的代碼示例,我們展示了如何創(chuàng)建和添加表格,更改單元格的樣式,需要的朋友可以參考下2025-12-12
Python解決asyncio文件描述符最大數(shù)量限制的問(wèn)題
這篇文章主要介紹了Python解決asyncio文件描述符最大數(shù)量限制的問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2024-06-06
Python將運(yùn)行結(jié)果導(dǎo)出為CSV格式的兩種常用方法
這篇文章主要給大家介紹了關(guān)于Python將運(yùn)行結(jié)果導(dǎo)出為CSV格式的兩種常用方法,Python生成(導(dǎo)出)csv文件其實(shí)很簡(jiǎn)單,我們一般可以用csv模塊或者pandas庫(kù)來(lái)實(shí)現(xiàn),需要的朋友可以參考下2023-07-07

