使用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
相关产品推荐
相关产品推荐

