python腳本將mysql數(shù)據(jù)寫入doris過程
python腳本將mysql數(shù)據(jù)寫入doris
flink-cdc連接的FE,數(shù)據(jù)寫入時正常的
但通過python腳本使用stream load方式直接連接FE,提示
NOT_AUTHORIZED]no valid Basic authorization
如果改用直接連接BE節(jié)點(diǎn),寫入又提示
[DATA_QUALITY_ERROR]too many filtered rows
這個就是搞技術(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 萬條
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 請求
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("錯誤信息:{}", 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列了。
看來問題出在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 "時柱",
`year_xun` CHAR(2) COMMENT "年柱旬",
`month_xun` CHAR(2) COMMENT "月柱旬",
`day_xun` CHAR(2) COMMENT "日柱旬",
`hour_xun` CHAR(2) COMMENT "時柱旬",
`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 "時干",
`hz` CHAR(1) COMMENT "時支",
`year_shu` TINYINT COMMENT "年柱數(shù)",
`month_shu` TINYINT COMMENT "月柱數(shù)",
`day_shu` TINYINT COMMENT "日柱數(shù)",
`hour_shu` TINYINT COMMENT "時柱數(shù)",
`yz_cs` CHAR(2) COMMENT "年長生",
`mz_cs` CHAR(2) COMMENT "月長生",
`dz_cs` CHAR(2) COMMENT "日長生",
`hz_cs` CHAR(2) COMMENT "時長生",
`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 "時干生克",
`hz_sk` CHAR(1) COMMENT "時支生克",
`year_ny` CHAR(3) COMMENT "年納音",
`month_ny` CHAR(3) COMMENT "月納音",
`day_ny` CHAR(3) COMMENT "日納音",
`hour_ny` CHAR(3) COMMENT "時納音",
`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 "正財數(shù)量",
`pc_cn` TINYINT COMMENT "偏財數(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 "劫財數(shù)量",
`ym_gj` CHAR(1) COMMENT "年月拱夾",
`md_gj` CHAR(1) COMMENT "月時拱夾",
`dh_gj` CHAR(1) COMMENT "日時拱夾"
)
ENGINE=OLAP
DUPLICATE KEY(`id`)
DISTRIBUTED BY HASH(`id`) BUCKETS 10
PROPERTIES (
"replication_num" = "1"
);
但實(shí)際寫入,報下面的錯誤,看來錯誤原因應(yīng)該是mysql的字節(jié)與doris的字節(jié)計算不一樣。
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í)行下面的語句,會提示
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 字符集,包括表情符號等),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 "時柱",
`year_xun` VARCHAR(6) COMMENT "年柱旬",
`month_xun` VARCHAR(6) COMMENT "月柱旬",
`day_xun` VARCHAR(6) COMMENT "日柱旬",
`hour_xun` VARCHAR(6) COMMENT "時柱旬",
`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 "時干",
`hz` VARCHAR(9) COMMENT "時支",
`year_shu` TINYINT COMMENT "年柱數(shù)",
`month_shu` TINYINT COMMENT "月柱數(shù)",
`day_shu` TINYINT COMMENT "日柱數(shù)",
`hour_shu` TINYINT COMMENT "時柱數(shù)",
`yz_cs` VARCHAR(6) COMMENT "年長生",
`mz_cs` VARCHAR(6) COMMENT "月長生",
`dz_cs` VARCHAR(6) COMMENT "日長生",
`hz_cs` VARCHAR(6) COMMENT "時長生",
`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 "時干生克",
`hz_sk` VARCHAR(9) COMMENT "時支生克",
`year_ny` VARCHAR(9) COMMENT "年納音",
`month_ny` VARCHAR(9) COMMENT "月納音",
`day_ny` VARCHAR(9) COMMENT "日納音",
`hour_ny` VARCHAR(9) COMMENT "時納音",
`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 "正財數(shù)量",
`pc_cn` TINYINT COMMENT "偏財數(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 "劫財數(shù)量",
`ym_gj` VARCHAR(9) COMMENT "年月拱夾",
`md_gj` VARCHAR(9) COMMENT "月時拱夾",
`dh_gj` VARCHAR(9) COMMENT "日時拱夾"
)
ENGINE=OLAP
UNIQUE KEY(`id`)
DISTRIBUTED BY HASH(`id`) BUCKETS 10
PROPERTIES (
"replication_num" = "1"
);
擴(kuò)展后,數(shù)據(jù)寫入正常了,于是我又驗(yàn)證了,反復(fù)執(zhí)行,看看有沒有問題。結(jié)果在doris中出現(xiàn)了兩條數(shù)據(jù)。id不是key,為什么會重復(fù)寫入呢?
因?yàn)樵贒oris中,“Duplicate Key”是一個冗余模型的特性。這個模型的數(shù)據(jù)完全按照導(dǎo)入文件中的數(shù)據(jù)進(jìn)行存儲,不會有任何聚合。即使兩行數(shù)據(jù)完全相同,也都會保留。
這個Duplicate模型針對日志是可以
但是針對我們的系統(tǒng)表是不合適的,這里就需要采用UNIQUE KEY

51萬數(shù)據(jù),寫入過程日志,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這個查詢比mysql快了24倍。那么為什么doris那么快呢?
總結(jié)
以上為個人經(jīng)驗(yàn),希望能給大家一個參考,也希望大家多多支持腳本之家。
相關(guān)文章
pandas中的DataFrame按指定順序輸出所有列的方法
下面小編就為大家分享一篇pandas中的DataFrame按指定順序輸出所有列的方法,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧2018-04-04
Django的ListView超詳細(xì)用法(含分頁paginate)
這篇文章主要介紹了Django的ListView超詳細(xì)用法(含分頁paginate),文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2020-05-05
Django2.1集成xadmin管理后臺所遇到的錯誤集錦(填坑)
這篇文章主要介紹了Django2.1集成xadmin管理后臺所遇到的錯誤集錦(填坑),小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧2018-12-12
python開發(fā)實(shí)例之Python的Twisted框架中Deferred對象的詳細(xì)用法與實(shí)例
這篇文章主要介紹了python開發(fā)實(shí)例之Python的Twisted框架中Deferred對象的詳細(xì)用法與實(shí)例,需要的朋友可以參考下2020-03-03
Python解決asyncio文件描述符最大數(shù)量限制的問題
這篇文章主要介紹了Python解決asyncio文件描述符最大數(shù)量限制的問題,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教2024-06-06
Python將運(yùn)行結(jié)果導(dǎo)出為CSV格式的兩種常用方法
這篇文章主要給大家介紹了關(guān)于Python將運(yùn)行結(jié)果導(dǎo)出為CSV格式的兩種常用方法,Python生成(導(dǎo)出)csv文件其實(shí)很簡單,我們一般可以用csv模塊或者pandas庫來實(shí)現(xiàn),需要的朋友可以參考下2023-07-07

