Bonobo使用咨询:提取步骤如何接入多数据源并添加至数据管道?
嘿,很高兴看到你已经摸透Bonobo的基础玩法了!针对你问的「接入两种不同数据源(比如两个网站)到数据管道」的问题,其实Bonobo的DAG(有向无环图)设计天生就支持这种多数据源场景,下面给你拆解两种实用的实现方式,附代码示例:
方法1:并行独立处理多数据源(最常用)
这种方式适合两个数据源的处理流程类似的场景——每个数据源对应一个独立的提取函数,然后各自接入管道的转换、加载环节,Bonobo会自动并行执行这些任务。
举个实际的例子,假设我们要从两个不同网站抓取内容:
import bonobo import requests from bs4 import BeautifulSoup # 提取网站A的数据 def extract_website_a(): response = requests.get("https://example.com/page-a") soup = BeautifulSoup(response.text, "html.parser") # 这里根据网站A的结构抓取数据,比如文章标题 for title_tag in soup.select(".article-title"): yield {"source": "网站A", "content": title_tag.get_text(strip=True)} # 提取网站B的数据 def extract_website_b(): response = requests.get("https://example.com/page-b") soup = BeautifulSoup(response.text, "html.parser") # 网站B的结构不同,抓取产品名称 for product_tag in soup.select(".product-name"): yield {"source": "网站B", "content": product_tag.get_text(strip=True)} # 统一的转换函数:清洗数据格式 def clean_data(item): item["content"] = item["content"].replace("\r\n", "").strip() return item # 加载函数:输出到控制台(实际可以替换成写入数据库/文件) def load_data(item): print(f"已加载数据:{item}") def get_graph(): graph = bonobo.Graph() # 给每个数据源添加一条完整的处理链 graph.add_chain(extract_website_a, clean_data, load_data) graph.add_chain(extract_website_b, clean_data, load_data) return graph if __name__ == "__main__": bonobo.run(get_graph())
运行这段代码,你会看到两个网站的数据被同时抓取、处理并输出——Bonobo会自动管理并行任务,不用额外配置线程/进程。
方法2:合并多数据源后统一处理
如果希望先把两个数据源的数据流合并,再进行后续处理,可以通过手动连接节点的方式实现:
def get_graph(): graph = bonobo.Graph() # 定义各个节点 extract_a = graph.add_node(extract_website_a) extract_b = graph.add_node(extract_website_b) clean = graph.add_node(clean_data) load = graph.add_node(load_data) # 连接节点:两个提取节点的输出都流向同一个转换节点 graph.add_edge(extract_a, clean) graph.add_edge(extract_b, clean) graph.add_edge(clean, load) return graph
这种方式和方法1的最终效果类似,但更适合需要灵活调整节点关系的复杂场景。
额外技巧:针对不同数据源做专属处理
如果两个数据源需要不同的转换逻辑,只要给每个数据源的链配置专属的转换函数即可:
# 网站A专属的转换函数 def transform_for_a(item): item["content"] = f"[网站A] {item['content']}" return item # 网站B专属的转换函数 def transform_for_b(item): item["content"] = f"[网站B] {item['content']}" return item def get_graph(): graph = bonobo.Graph() graph.add_chain(extract_website_a, transform_for_a, load_data) graph.add_chain(extract_website_b, transform_for_b, load_data) return graph
核心思路就是每个数据源对应一个独立的提取函数,然后根据你的需求,把这些提取节点接入到管道的不同位置——Bonobo的DAG模型会帮你处理剩下的流程调度工作。
内容的提问来源于stack exchange,提问作者Volatil3
相关产品推荐
相关产品推荐

