【问题标题】:Converting large text/pgn files to JSON in Spark在 Spark 中将大型文本/pgn 文件转换为 JSON
【发布时间】:2021-07-05 05:09:52
【问题描述】:

我需要将 PGN 文件转换为 JSON,然后我可以使用 Spark 将它们转换为 Spark DataFrames 并最终创建一个图表。我已经编写了一个 python 脚本来使用 Pandas 将它们解析为 DataFrame,但它太慢了(170k 游戏大约需要 56 分钟(原始估计为 30 分钟,但在配置文件后我估计为 56 分钟))。我还尝试使用这个 repo:https://github.com/JonathanCauchi/PGN-to-JSON-Parser,它给了我 JSON 文件,但 170k 游戏花了 69 分钟。

我可以将 PGN 扩展名更改为 .txt,它的工作方式似乎完全相同,所以我认为对 .txt 到 JSON 的支持更多,但我不确定。

我认为 Spark 会比“普通”Python 更快,但我不知道如何进行转换。下面是一个示例。虽然有 20 亿个游戏,所以我目前的方法都不起作用,因为如果我要使用 PGN-to-JSON-Parser 需要将近 2 年的时间。理想情况下,一个 .txt 到 Spark DataFrames 并完全忽略 JSON 将是理想的。

[Event "Rated Classical game"]
[Site "https://lichess.org/j1dkb5dw"]
[White "BFG9k"]
[Black "mamalak"]
[Result "1-0"]
[UTCDate "2012.12.31"]
[UTCTime "23:01:03"]
[WhiteElo "1639"]
[BlackElo "1403"]
[WhiteRatingDiff "+5"]
[BlackRatingDiff "-8"]
[ECO "C00"]
[Opening "French Defense: Normal Variation"]
[TimeControl "600+8"]
[Termination "Normal"]

1. e4 e6 2. d4 b6 3. a3 Bb7 4. Nc3 Nh6 5. Bxh6 gxh6 6. Be2 Qg5 7. Bg4 h5 8. Nf3 Qg6 9. Nh4 Qg5 10. Bxh5 Qxh4 11. Qf3 Kd8 12. Qxf7 Nc6 13. Qe8# 1-0

编辑:添加了 20k 游戏的配置文件。

ncalls  tottime  percall  cumtime  percall filename:lineno(function)
        1   26.828   26.828  395.848  395.848 /Users/danieljones/Documents – Daniel’s iMac/GitHub/ST446Project/ParsePGN.py:11(parse_pgn)
    20000    0.798    0.000  289.203    0.014 /Users/danieljones/opt/anaconda3/envs/LSE/lib/python3.6/site-packages/pandas/core/frame.py:7614(append)
    20000    0.098    0.000  199.489    0.010 /Users/danieljones/opt/anaconda3/envs/LSE/lib/python3.6/site-packages/pandas/core/reshape/concat.py:70(concat)
    20000    0.480    0.000  126.548    0.006 /Users/danieljones/opt/anaconda3/envs/LSE/lib/python3.6/site-packages/pandas/core/reshape/concat.py:295(__init__)
   100002    0.212    0.000  122.178    0.001 /Users/danieljones/opt/anaconda3/envs/LSE/lib/python3.6/site-packages/pandas/core/generic.py:5199(_protect_consolidate)
    80002    0.076    0.000  122.177    0.002 /Users/danieljones/opt/anaconda3/envs/LSE/lib/python3.6/site-packages/pandas/core/generic.py:5210(_consolidate_inplace)
    40000    0.079    0.000  122.063    0.003 /Users/danieljones/opt/anaconda3/envs/LSE/lib/python3.6/site-packages/pandas/core/generic.py:5218(_consolidate)
    80002    0.170    0.000  121.830    0.002 /Users/danieljones/opt/anaconda3/envs/LSE/lib/python3.6/site-packages/pandas/core/generic.py:5213(f)
   100001    0.223    0.000   99.829    0.001 /Users/danieljones/opt/anaconda3/envs/LSE/lib/python3.6/site-packages/pandas/core/internals/managers.py:986(_consolidate_inplace)
    59999    0.451    0.000   96.718    0.002 /Users/danieljones/opt/anaconda3/envs/LSE/lib/python3.6/site-packages/pandas/core/internals/managers.py:1898(_consolidate)
    80002    0.138    0.000   96.599    0.001 /Users/danieljones/opt/anaconda3/envs/LSE/lib/python3.6/site-packages/pandas/core/internals/managers.py:970(consolidate)
    79999   52.602    0.001   91.913    0.001 /Users/danieljones/opt/anaconda3/envs/LSE/lib/python3.6/site-packages/pandas/core/internals/managers.py:1915(_merge_blocks)
    20000    7.432    0.000   79.741    0.004 /Users/danieljones/opt/anaconda3/envs/LSE/lib/python3.6/site-packages/chess/pgn.py:1323(read_game)
    20000    0.361    0.000   72.843    0.004 /Users/danieljones/opt/anaconda3/envs/LSE/lib/python3.6/site-packages/pandas/core/reshape/concat.py:456(get_result)

我不确定“cumtime”是否是排序的最佳列,但似乎附加步骤需要很多时间。

这是我的脚本:

def parse_pgn(pgn):
    games = []
    i = 0
    edges_df = pd.DataFrame(columns=["Event", "Round", "WhitePlayer", "BlackPlayer", "Result", "BlackElo",
                                     "Opening", "TimeControl", "Date", "Time", "WhiteElo"])

    while i < 20000:
        first_game = chess.pgn.read_game(pgn)

        if first_game is not None:
            Event = first_game.headers["Event"]
            Round = first_game.headers["Round"]
            White_player = first_game.headers["White"]
            Black_player = first_game.headers["Black"]
            Result = first_game.headers["Result"]  # Add condition to split this
            if Result == "1-0":
                Result = White_player
            elif Result == "0-0":
                Result = "Draw"
            else:
                Result = Black_player
            BlackELO = first_game.headers["BlackElo"]
            Opening = first_game.headers["Opening"]
            TimeControl = first_game.headers["TimeControl"]
            UTCDate = first_game.headers["UTCDate"]
            UTCTime = first_game.headers["UTCTime"]
            WhiteELO = first_game.headers["WhiteElo"]
            edges_df = edges_df.append({"Event": Event,
                                                "Round": Round,
                                                "WhitePlayer": White_player,
                                                "BlackPlayer": Black_player,
                                                "Result": Result,
                                                "BlackElo": BlackELO,
                                                "Opening": Opening,
                                                "TimeControl": TimeControl,
                                                "Date": UTCDate,
                                                "Time": UTCTime,
                                                "White": WhiteELO,
                                                }, ignore_index=True)
            games.append(first_game)
            i += 1
        else:
            pass

    return edges_df

编辑 2:将 append 方法更改为字典。 20k 现在需要 78 秒。很多花时间的方法似乎都来自chess 包,例如检查合法动作、阅读棋盘布局。这些对我的最终目标都不重要,所以我想知道我是否可以不再使用这个包,而是自己将文件拆分为单独的游戏,也许在[Event,因为这是每个不同游戏的开始。

【问题讨论】:

  • 已经编写了一个 python 脚本来使用 Pandas 将它们解析为 DataFrame,但它太慢了(170k 游戏大约需要 30 分钟)您是否尝试过分析或优化您的代码?
  • 好主意,谢谢。我已将此添加为编辑。

标签: python json apache-spark pyspark distributed-computing


【解决方案1】:

如果您想缩短运行时间,请不要将.appendpandas.DataFrame 循环,您可以阅读更多关于此here 的信息。您可能首先将您的 dicts 存储在 iterable 中,然后从中创建 pandas.DataFrame。我会使用collections.deque(来自collections 内置模块),因为它旨在运动高速.append,让我们比较这些不同的方式

import collections
import pandas as pd
def func1():
    df = pd.DataFrame(columns=['x','y','z'])
    for i in range(1000):
        df = df.append({'x':i,'y':i*10,'z':i*100}, ignore_index=True)
    return df
def func2():
    store = collections.deque()
    for i in range(1000):
        store.append({'x':i,'y':i*10,'z':i*100})
    df = pd.DataFrame(store, columns=['x','y','z'])
    return df

这些函数产生相等的pandas.DataFrames,我使用内置模块timeit按照以下方式比较它们

import timeit
print(timeit.timeit('func1()',number=10,globals={'func1':func1}))
print(timeit.timeit('func2()',number=10,globals={'func2':func2}))

得到以下结果

18.370871236000312
0.02604325699940091

速度快了 500 多倍。自然您的里程可能会有所不同,但我建议尝试这种优化。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-07-14
    • 2012-06-17
    • 1970-01-01
    • 2013-09-16
    • 1970-01-01
    • 2020-06-12
    • 1970-01-01
    相关资源
    最近更新 更多