最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

python腳本將mysql數(shù)據(jù)寫入doris過程

 更新時間:2026年02月25日 09:48:11   作者:warrah  
文章描述了在將MySQL數(shù)據(jù)寫入Doris時遇到的權(quán)限和數(shù)據(jù)質(zhì)量問題,包括使用Python腳本和Streamload方式的異常,以及通過Flink-CDC進(jìn)行數(shù)據(jù)同步時遇到的字符集和主鍵問題,作者通過分析和調(diào)整,最終解決了數(shù)據(jù)寫入問題,并對比了Doris和MySQL在性能上的差異

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)文章

最新評論

邵武市| 塘沽区| 岫岩| 襄樊市| 北川| 浙江省| 汝阳县| 子洲县| 内丘县| 陆丰市| 汉源县| 甘德县| 东乌珠穆沁旗| 尼玛县| 新源县| 天全县| 泸定县| 合阳县| 巴林左旗| 镇雄县| 海林市| 弥勒县| 久治县| 信宜市| 临猗县| 巴彦淖尔市| 织金县| 天门市| 汕尾市| 云林县| 宜兰县| 息烽县| 丹东市| 安福县| 利辛县| 曲麻莱县| 洪湖市| 扎赉特旗| 广西| 尼木县| 吉林省|