Skip to content
farfarfunPublic

About

阿里天池 531800 赛题参赛代码:基于 Flink AI Flow 编排的流批一体推荐工作流,含 Kafka 数据源、AutoEncoder 训练与在线预测、Proxima 向量索引构建与相似检索

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Latest commit

 

History

33 Commits

Folders and files

Repository files navigation

funfight

一次性的比赛代码存档,不是通用工具/库,目前已废弃、未维护。

仓库里保存的是作者参加阿里云天池比赛(比赛编号 531800)时提交的方案代码:基于阿里开源的 AI Flow 框架编写的一套批流一体机器学习工作流,用 TensorFlow 训练自编码器(Autoencoder)、用 Proxima 构建向量索引、再用 Flink 做在线预测/检索。除了这一份比赛代码外,没有其他功能。

说明:目前 PyPI 上查不到 funfight 已发布的版本(返回 Not Found)。本仓库只是一次性比赛代码存档,不建议 pip install。

目录结构 / 代码做了什么

src/funfight/
└── tianchi/
    └── t531800/
        ├── step1.py            # 准备环境:从蓝奏云下载数据集,安装 apache-flink / kafka-python,下载 Flink/Kafka 安装包
        ├── kafka_source.py     # Source 类:监听 AIFlow 通知,创建/清理 Kafka topic,把测试集逐行发到 Kafka 作为在线推理输入
        ├── ai_flow_master.py   # 启动 AIFlowMaster(读取同目录 master.yaml)
        └── package/python_codes/
            ├── tianchi_main.py       # 定义完整 workflow:register_example/register_model 注册数据与模型 →
            │                         #   训练自编码器(TrainAutoEncoder) → cluster_serving 上线 →
            │                         #   Flink 任务构建向量索引(BuildIndexExecutor) → 预测(PredictAutoEncoder) →
            │                         #   向量检索(SearchExecutor/SearchExecutor3) → 写出结果(WriteSecondResult)
            ├── python_job_executor.py  # ReadCsvExample / TrainAutoEncoder:读 CSV、训练一个简单的 Dense 自编码器并保存
            ├── tianchi_executor.py     # Flink 侧的 Executor:读数据、预测、拼接历史、写结果等
            ├── proxima_executor.py     # BuildIndexExecutor / SearchExecutor:调用 Proxima 建索引、做近邻检索
            └── data_type.py            # FloatDataType / DoubleDataType:在 Proxima 类型和 Flink 类型之间转换

example/531800.py 是空文件;src/funfight/__init__.py、src/funfight/tianchi/__init__.py、t531800/__init__.py 里只有说明性的模块 docstring,没有对外暴露任何 API。

依赖

pyproject.toml 声明的直接依赖是本组织的 fundrive[lanzou](下载数据集用)、farlog(日志)、funshell(执行 shell 命令)。

除此之外,src/funfight/tianchi/t531800/ 下的比赛方案代码还依赖比赛当年的特定环境,这些依赖都没有写进 pyproject.toml,pip install funfight 装不出一个能跑的环境(2026-10 复核):

导入名 公开 PyPI 包 现状
ai_flow ai-flow 0.1.0 存在,但声明 requires-python >=3.7,<3.8,装不进本仓库要求的 Python ≥3.10
flink_ai_flow、python_ai_flow 无 未发布到 PyPI,只能从 flink-extended/ai-flow 源码或比赛提供的包里取
pyproxima2 无 阿里内部向量检索库,未公开发布
zoo.serving.client analytics-zoo 存在,但同样是面向老版本 Python/Spark 的历史包
pyflink apache-flink 需锁在 1.11.0,与现代版本 API 不兼容
kafka-python、tensorflow、pandas、numpy、PyYAML 有 版本需与当年环境匹配

step1.py 里还有 pip install apache-flink==1.11.0、下载 Flink 1.11.0 / Kafka 2.3.0 安装包等步骤,整套环境无法在现代机器上直接复现,需要按天池比赛当年的环境手动搭建。另外 step1() 用的两个蓝奏云转存链接(wws.lanzous.com)已经失效,数据集请从赛题页面重新下载。

安装

本仓库未发布到 PyPI(pip install funfight 会报 Not Found),也不建议安装使用;获取代码请直接克隆仓库:

git clone https://github.com/farfarfun/funfight.git

如果确实需要把 funfight 当作本地包安装(例如只是为了让 import funfight 不报错),可以用 uv/pip 从源码安装:

uv pip install -e .
# 或
pip install -e .

最小可运行示例

除比赛方案代码外,funfight 包本身不提供任何公共 API,唯一能验证安装是否成功的方式是导入空包:

python3 -c "import funfight; print(funfight.__file__)"

使用(服务/流程入口)

没有命令行入口,也没有可直接调用的公共函数。这些流程依赖比赛当年的历史环境(Flink 1.11.0、Kafka 2.3.0、阿里内部的 AI Flow / Proxima 服务等),不是可以用 scripts/setup.sh 统一 start/stop 的常规服务,因此仍按当年的手工步骤说明,不提供伪装成「一键启停」的脚本:

  1. 执行 src/funfight/tianchi/t531800/step1.py 里的 step1()/step2()/step3(),依次下载数据集、安装 Flink/Kafka 历史版本、准备 ai_flow 环境;
  2. 参考 src/funfight/tianchi/t531800/ai_flow_master.py 启动 AIFlowMaster(读取同目录 master.yaml);
  3. 按 src/funfight/tianchi/t531800/README.md 配置好 PYTHONPATH/ENV_HOME/TASK_ID/FLINK_HOME 等环境变量、并把 source.yaml 的 dataset_uri 指向 second_test_data.csv 后,启动 src/funfight/tianchi/t531800/kafka_source.py(Kafka Source);
  4. 最后运行 src/funfight/tianchi/t531800/package/python_codes/tianchi_main.py 提交 workflow。

子目录 README 里列齐了代码实际读取的三个数据文件(train_data.csv、first_test_data.csv、second_test_data.csv)、每个环境变量的读取位置,以及产物的落盘路径。


关于 farfarfun

farfarfun 是一个专注于实用工具库的开源组织, 涵盖云存储、数据处理、AI、多媒体与开发工具链等方向。

本项目基于 MIT 协议开源。

About

阿里天池 531800 赛题参赛代码:基于 Flink AI Flow 编排的流批一体推荐工作流,含 Kafka 数据源、AutoEncoder 训练与在线预测、Proxima 向量索引构建与相似检索

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Used by

Contributors

Languages