You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用HuggingFace Datasets调用Apache FlinkRunner时出现Client Cancelled异常的问题求助

解决Hugging Face Datasets下载维基百科(ceb/es/ru)时的Beam超时问题

看起来你遇到的CANCELLED: client cancelled错误大概率是Colab会话资源限制(内存/磁盘/超时)或者Flink Runner在Colab环境下的兼容性问题导致的——毕竟Colab的免费版会话通常会在12小时左右断开,而且分布式Beam运行在单实例环境里很容易出问题。结合你时间紧迫的实习项目需求,给你几个优先级从高到低的解决方案:

1. 优先用流式加载+分批处理,跳过Beam的分布式依赖

直接放弃Flink Runner,改用Hugging Face Datasets的流式加载功能,分批处理并保存,这样能避免一次性加载整个大语料导致的内存溢出和会话超时。修改后的代码如下:

!pip install datasets mwparserfromhell
import os
from datasets import load_dataset, Dataset
from google.colab import drive

drive_dir = os.path.join(os.getcwd(), 'drive')
drive.mount(drive_dir)

lang = 'ru' # 替换为'ceb'或'es'
lang_dir = os.path.join(drive_dir, 'path/to/training/dir', lang)

if not os.path.exists(lang_dir):
    # 启用流式加载,避免一次性占用大量内存
    ds_stream = load_dataset('wikipedia', '20220301.' + lang, split='train', streaming=True)
    
    # 分批处理并收集数据
    batch_size = 1000
    batches = []
    for idx, batch in enumerate(ds_stream.iter(batch_size=batch_size)):
        batches.append(batch)
        if idx % 10 == 0:
            print(f"已处理 {idx*batch_size} 条数据")
    
    # 合并批次为完整数据集并保存
    full_ds = Dataset.from_dict({
        key: sum([b[key] for b in batches], []) 
        for key in batches[0].keys()
    })
    full_ds.save_to_disk(lang_dir)

这种方法不需要Apache Beam,完全在单进程里处理,适配Colab的资源限制,而且不会触发分布式相关的错误。

2. 直接用预处理好的现成数据集(最快解决方案)

如果时间实在紧张,不用自己下载预处理,Hugging Face Hub上有很多用户分享的已经处理好的维基百科语料:

  • 俄语:可以用IlyaGusev/ru_wikipedia,已经处理成纯文本,直接加载即可
  • 西班牙语:比如datificate/wikipedia-es,包含不同版本的预处理文本
  • 宿务语:facebook/ceb-wikipedia是官方预训练用的预处理版本

加载示例:

from datasets import load_dataset
# 加载俄语预处理数据集
ds = load_dataset("IlyaGusev/ru_wikipedia", split="train")
ds.save_to_disk(lang_dir)

这些数据集都是已经完成清洗、提取文本的版本,直接就能用,完全省去下载和预处理的时间。

3. 优化Colab环境,避免会话中断

如果一定要自己下载原始数据集,可以通过以下方式延长Colab会话:

  • 切换到Colab Pro/Pro+:获得更大的内存(最高32GB)和更长的会话时间(最多24小时),避免因为资源不足被强制终止
  • 保持会话活跃:在浏览器控制台运行以下脚本,每隔1分钟点击一次连接按钮,防止Colab检测到无交互而断开:
    function ClickConnect(){
        console.log("保持会话活跃");
        document.querySelector("colab-connect-button").click()
    }
    setInterval(ClickConnect, 60000)
    
  • 清理磁盘空间:运行!rm -rf /content/sample_data删除默认示例数据,确保有足够空间存放语料

4. 调整Beam配置(仅当必须用Beam时)

如果坚持要使用Beam处理,建议换成DirectRunner(单进程运行),并限制资源使用,避免分布式相关的错误:

!pip install datasets apache_beam mwparserfromhell
import os
from datasets import load_dataset
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions
from google.colab import drive

drive_dir = os.path.join(os.getcwd(), 'drive')
drive.mount(drive_dir)

lang = 'ru'
lang_dir = os.path.join(drive_dir, 'path/to/training/dir', lang)

if not os.path.exists(lang_dir):
    # 设置Beam选项,限制单worker和内存
    beam_options = PipelineOptions()
    worker_opts = beam_options.view_as(beam.options.pipeline_options.WorkerOptions)
    worker_opts.num_workers = 1
    worker_opts.autoscaling_algorithm = 'NONE'
    resource_opts = beam_options.view_as(beam.options.pipeline_options.ResourceOptions)
    resource_opts.memory_mb = 8192  # 根据Colab可用内存调整
    
    # 使用DirectRunner替代Flink
    x = load_dataset(
        'wikipedia', '20220301.' + lang, 
        beam_runner='DirectRunner',
        beam_options=beam_options,
        split='train'
    )
    x.save_to_disk(lang_dir)

以上几种方法里,我最推荐前两种——流式处理或者直接用现成数据集,能最快解决你的问题,适配实习项目的时间要求。

内容的提问来源于stack exchange,提问作者Dan DeGenaro

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.28 23:59:05