title: "TRON 区块链数据采集:基于 Scrapy 与 Flask 的链上数据服务构建" slug: "tron-blockchain-crawler-project" date: 2026-06-23 tags: [TRON, Scrapy, Flask, MongoDB, 区块链, 爬虫]
前言
TRON(波场)作为全球知名的公有区块链平台,其每秒数千笔的交易吞吐量为链上数据分析带来了持续的挑战。无论是做地址归集、交易图谱分析,还是构建实时大屏看板,第一步永远是高效、稳定地采集区块与交易数据。
近期我在 GitHub 上发现了一个很有意思的开源项目——tron_address,它巧妙地将 Scrapy 爬虫引擎与 Flask Web 框架结合起来,围绕 TRON Grid API 构建了一套完整的区块数据采集与控制服务。本文将深入剖析它的架构设计与实现细节,希望能为正在做区块链数据基建的开发者提供一些参考。
一、项目总体架构
在展开代码之前,我们先从宏观上看一下这个项目的数据流:
TRON Grid API (api.trongrid.io)
│
▼
┌─────────────────┐ ┌──────────────────┐ ┌──────────────────┐
│ Scrapy Spider │────▶│ Pipeline 管道 │────▶│ MongoDB │
│ (tron_spider) │ │ (MongoDBPipeline) │ │ (区块 + 交易) │
└──────┬──────────┘ └──────────────────┘ └──────────────────┘
│
│ (通过线程启动/停止)
▼
┌─────────────────┐ ┌──────────────────┐
│ Flask 控制服务 │────▶│ Monitor/Controller│
│ (app.py/run.py) │ │ (实时监控) │
└─────────────────┘ └──────────────────┘
整个系统分为两大模块:
- 数据采集层——基于 Scrapy 框架的
tron_spider.py,负责从 TRON Grid API 拉取区块和交易数据,经过 Item Pipeline 写入 MongoDB。 - 控制与监控层——基于 Flask 的 HTTP 服务,提供爬虫启停 API、运行状态查询,以及实时的区块更新监控。
此外,configuration.py 负责配置管理,key.json 提供了 API Key 轮换机制。这种"采集 + 控制"分离的设计,让系统既能独立运行爬虫,又能通过 API 灵活编排任务。
二、Scrapy 爬虫设计:tron_spider 深度解析
2.1 爬虫初始化与启动策略
tron_spider.py 是整个系统的数据引擎。它的启动方式非常巧妙——先请求 getnowblock 接口获取当前最新区块号,再以此为依据切分批量任务:
class TronSpider(scrapy.Spider):
name = 'tron_spider'
def __init__(self, *args, **kwargs):
super(TronSpider, self).__init__(*args, **kwargs)
# 初始化 MongoDB 连接和集合
self.client = MongoClient(self.settings.get('MONGODB_URI'))
self.db = self.client[self.settings.get('MONGODB_DATABASE')]
self.start_block = kwargs.get('start_block', 1078_5503)
self.batch_size = kwargs.get('batch_size', 200000)
self.max_concurrent_requests = kwargs.get('max_concurrent_requests', 500)
这里有几个值得关注的细节:
- 可配置的起始区块号:通过
start_block参数指定,支持从任意高度开始采集,而不是必须从创世块开始。 - 批量处理:
batch_size = 200000意味着每次只请求 20 万个区块的数据,避免一次性产生海量请求。 - MongoDB 连接在爬虫内直接管理:
spider_closed信号确保爬虫关闭时优雅断开数据库连接。
2.2 两阶段请求流程
爬虫的核心请求策略分为两个阶段:
阶段一:获取最新高度
def parse(self, response):
data = json.loads(response.text)
latest_block_number = data["block_header"]["raw_data"]["number"]
if self.is_first_run:
self.is_first_run = False
self.current_batch_end = min(
self.current_batch_start + self.batch_size,
latest_block_number + 1
)
# 生成当前批次的区块请求
for block_num in range(self.current_batch_start, self.current_batch_end):
url = f"https://api.trongrid.io/wallet/getblockbynum?num={block_num}"
yield Request(url=url, callback=self.parse_block, meta={'block_num': block_num})
首次运行先确定"我要抓取的范围",然后批量下发请求。
阶段二:逐块解析与自动续批
def parse_block(self, response):
block_num = response.meta['block_num']
data = json.loads(response.text)
# 提取区块 Item
yield self.process_block_data(data)
# 提取交易 Item
for tx in data.get("transactions", []):
yield self.process_transaction_data(tx, block_num, data.get("blockID"))
# 当前批次完成后自动查询最新高度,开启下一批次
if block_num == self.current_batch_end - 1:
yield Request(
url="https://api.trongrid.io/wallet/getnowblock",
callback=self.start_new_batch,
)
每一批的最后一块处理完后,自动触发 start_new_batch,再次查询最新高度,决定是否需要继续。这种链式调度让爬虫具备了"永远追赶最新区块"的能力,类似于一个永动机。
2.3 区块与交易的数据提取
process_block_data 和 process_transaction_data 两个方法负责从原始 JSON 中提取有用字段,同时做去重检查——如果区块或交易已经存在于数据库,则跳过,避免重复写入:
def process_block_data(self, block_data):
block_number = block_header.get("number")
if self.blocks_collection.find_one({"number": block_number}):
return None # 已存在,跳过
block_item = TronBlockItem()
block_item['blockID'] = block_id
block_item['number'] = block_number
block_item['timestamp'] = block_header.get("timestamp")
block_item['witness_address'] = block_header.get("witness_address")
block_item['created_at'] = datetime.now()
return block_item
这种"查重再写"的策略虽然多了一次数据库查询,但对于区块链这种不可变数据源来说,能有效保证数据的幂等性。
2.4 Item 定义与数据流转
项目在 items.py 中定义了三个 Scrapy Item 类,分别对应核心数据实体:
class TronBlockItem(scrapy.Item):
blockID = scrapy.Field() # 区块哈希
number = scrapy.Field() # 区块高度
timestamp = scrapy.Field() # 出块时间戳
txTrieRoot = scrapy.Field() # 交易树根哈希
parentHash = scrapy.Field() # 父块哈希
witness_address = scrapy.Field() # 出块节点地址
version = scrapy.Field() # 协议版本号
created_at = scrapy.Field() # 入库时间
class TronTransactionItem(scrapy.Item):
txID = scrapy.Field() # 交易哈希
block_number = scrapy.Field() # 所属区块高度
block_id = scrapy.Field() # 所属区块哈希
raw_data = scrapy.Field() # 交易原始数据
signature = scrapy.Field() # 签名数组
created_at = scrapy.Field() # 入库时间
class TronAddressItem(scrapy.Item):
address = scrapy.Field() # TRON 地址
block_num = scrapy.Field() # 首次出现区块
TronAddressItem 是另一个爬虫 tron_addr.py 使用的 Item,专门用于地址提取场景。它从每个区块的 witness_address 以及交易合约中的 owner_address、to_address、contract_address 字段提取所有关联地址,并通过 hex_to_base58 方法将 TRON 的 Hex 格式地址转换为用户友好的 Base58 格式(以 T 开头的地址)。
2.5 Pipeline 管道:MongoDB 写入与地址采集
Scrapy 的 Pipeline 机制是数据流转的最后环节。项目实现了三个 Pipeline:
MongoDBPipeline 是核心存储管道,负责将 Item 写入 MongoDB:
class MongoDBPipeline:
def open_spider(self, spider):
self.client = pymongo.MongoClient(self.mongodb_uri)
self.db = self.client[self.mongodb_db]
# 启动时自动创建唯一索引
self.db.blocks.create_index([("number", pymongo.ASCENDING)], unique=True)
self.db.transactions.create_index([("txID", pymongo.ASCENDING)], unique=True)
def process_item(self, item, spider):
if isinstance(item, TronBlockItem):
collection = self.db.blocks
elif isinstance(item, TronTransactionItem):
collection = self.db.transactions
item_dict = ItemAdapter(item).asdict()
collection.insert_one(item_dict)
return item
这里使用了 ItemAdapter 来统一处理不同 Item 类型,是 Scrapy 2.x 之后推荐的做法。Pipeline 在 open_spider 阶段建立连接并创建索引,在 close_spider 阶段关闭连接,生命周期与爬虫完全一致。
TronScannerPipeline 则是另一个轻量管道,专注于地址去重计数:
class TronScannerPipeline:
def __init__(self):
self.addresses = set()
def process_item(self, item, spider):
self.addresses.add(item['address'])
spider.crawler.stats.set_value('addresses_found', len(self.addresses))
return item
它利用 Python 的 set 天然去重特性,在内存中维护已发现地址集合,并通过 Scrapy 的 Stats 收集器暴露统计指标,供外部 API 查询。
三、TRON Grid API 对接实践
3.1 接口选择
项目主要调用了 TRON Grid 的两个核心 REST 接口:
| 接口 | 用途 | 请求方式 |
|---|---|---|
/wallet/getnowblock |
获取当前最新区块 | GET |
/wallet/getblockbynum?num={number} |
按编号获取指定区块详情 | GET |
这种设计非常务实——先获取"指针"(最新高度),再通过编号循环拉取详情。相比一次性拉取全量数据,这种方式对服务端更友好,也更容易实现断点续传。
3.2 API Key 轮换机制
TRON Grid API 对免费用户有频率限制,项目通过 configuration.py 实现了轻量级的 Key 轮换:
class Config:
def __init__(self):
self.__conf = read_json_file(CONFIG_PATH)
self.__conf["api_keys"] = load_api_keys()
def random_key(self):
keys = self.__conf.get("api_keys", [])
return random.choice(keys) if keys else None
- Key 来源有两种方式:环境变量
TRON_API_KEYS(逗号分隔)或key.json文件。 - 每次请求从 Key 列表中 随机选取,天然分散了单个 Key 的请求压力。
- 在
controller.py中,每次发起 HTTP 请求时都会将 Key 注入请求头:
async def async_get_data(self, url: str) -> Optional[Dict]:
headers = self.config.get("headers").copy()
headers["TRON-PRO-API-KEY"] = self.config.random_key()
async with session.get(url, headers=headers) as response:
return await response.json()
3.3 重试与容错
无论是 Scrapy 爬虫还是 controller 中的异步请求,都内置了重试机制。以 controller 为例:
async def get_block_data(self, d: int = 0) -> Optional[Dict]:
try:
response = await self.async_get_data(GETNOWBLOCK)
if response:
return response
raise Exception("空响应")
except Exception as e:
if d < self.config.get("RETRY_TIMES"):
await asyncio.sleep(self.config.get("RETRY_DELAY"))
return await self.get_block_data(d + 1)
重试次数和延迟间隔均可通过 config.json 配置,默认为 3 次重试、每次间隔 3 秒,充分适配链上 API 不稳定的场景。
3.4 地址 Hex 转 Base58
TRON 区块链的地址在链上以 Hex 格式存储(以 41 开头),而在钱包和浏览器中通常以 Base58 格式展示(以 T 开头)。项目在 controller.py 中实现了完整的转换逻辑:
@staticmethod
def hex_to_base58(hex_address: str) -> Optional[str]:
if not hex_address or not hex_address.startswith("41"):
return None
address_bytes = bytes.fromhex(hex_address)
hash0 = hashlib.sha256(address_bytes).digest()
hash1 = hashlib.sha256(hash0).digest()
checksum = hash1[:4]
return base58.b58encode(address_bytes + checksum).decode()
这个方法的实现遵循了 Base58Check 编码规范:对地址字节做双重 SHA256 哈希,取前 4 字节作为校验和,拼接后通过 Base58 编码输出。这样做不仅让地址更短、易于人工辨认,还能在校验和错误时提示用户输入有误。
在 tron_addr.py 爬虫中,地址提取的覆盖面更广——不仅提取出块地址 witness_address,还遍历每笔交易中所有合约类型的 owner_address、to_address 和 contract_address,做到不遗漏任何关联地址。
四、Flask 控制服务设计
项目提供了 两套 Flask 控制服务,分别解决不同的使用场景:
4.1 app.py —— 配置管理与实时监控
app.py 是一个轻量级 API 服务,路由设计清晰:
| 端点 | 方法 | 功能 |
|---|---|---|
/api/config |
GET | 获取所有配置项 |
/api/config/<key> |
GET | 获取指定配置项 |
/api/config/<key> |
PUT | 修改指定配置项 |
/api/get_cur_block |
GET | 获取当前最新区块数据 |
/api/get_last_time |
GET | 获取上次数据更新时间 |
其中,Monitor 类继承自 TronAddress,启动一个后台线程持续轮询最新区块,通过 MD5 哈希比对判断区块是否变化:
class Monitor(TronAddress):
def worker(self):
while True:
new_block_data = self.get_block_data_sync()
if new_block_data:
new_hash = self.get_data_hash(new_block_data)
if new_hash != self.cur_hash:
self.update(new_block_data)
self.cur_hash = new_hash
time.sleep(self.config.get("interval", 5))
当检测到新区块时,update 方法会增量计算地址交易次数变化,并通过 MongoDB 的 bulk_write + $inc 原子操作完成更新:
def submit_to_db(self, addresses):
bulk_ops = [
UpdateOne(
{'address': addr},
{'$inc': {'transaction_count': count}},
upsert=True
)
for addr, count in addresses.items()
]
collection.bulk_write(bulk_ops)
这种设计非常适合地址活跃度统计场景——不需要冗余存储每一笔交易,只需维护每个地址的累计交易次数。
4.2 run.py —— 爬虫生命周期管理
run.py 提供的是爬虫控制 API,面向需要远程管理 Scrapy 任务的场景:
@app.route('/start_crawler', methods=['POST'])
def start_crawler():
scrapy_thread = Thread(target=run_scrapy)
scrapy_thread.daemon = True
scrapy_thread.start()
return jsonify({'status': 'success', 'message': '爬虫启动成功'})
@app.route('/crawler_status', methods=['GET'])
def crawler_status():
if scrapy_thread and scrapy_thread.is_alive():
return jsonify({'status': 'running'})
return jsonify({'status': 'stopped'})
它通过守护线程来运行 Scrapy 的 CrawlerProcess,使得 Flask 服务本身不阻塞,同时又能通过 HTTP 接口随时查看爬虫运行状态。对于 Docker 容器化部署或 Kubernetes 集群来说,这样的 API 接口非常便于集成到运维监控系统中。
五、MongoDB 存储方案
5.1 集合设计
项目使用了 MongoDB 作为持久化存储,主要包含以下集合:
blocks 集合(区块数据):
{
"blockID": "000000000...", // 区块哈希
"number": 10785503, // 区块高度
"timestamp": 1700000000, // 出块时间戳
"txTrieRoot": "0x...", // 交易树根
"parentHash": "0x...", // 父区块哈希
"witness_address": "T...", // 出块地址
"version": 32, // 版本号
"created_at": ISODate(...) // 记录创建时间
}
transactions 集合(交易数据):
{
"txID": "0x...", // 交易哈希
"block_number": 10785503, // 所属区块高度
"block_id": "000000000...", // 所属区块哈希
"raw_data": { ... }, // 交易原始数据
"signature": [...], // 签名列表
"created_at": ISODate(...) // 记录创建时间
}
5.2 索引设计
在 MongoDBPipeline.open_spider 中,爬虫启动时会自动创建唯一索引:
self.db.blocks.create_index([("number", pymongo.ASCENDING)], unique=True)
self.db.transactions.create_index([("txID", pymongo.ASCENDING)], unique=True)
blocks.number唯一索引:保证同一个区块高度不会被重复存储。transactions.txID唯一索引:保证交易哈希全局唯一,天然防重。
5.3 数据迁移与批量导入
项目在 test/ 目录下提供了丰富的数据工具。其中 migration_morgodb.py 实现了跨环境的数据迁移功能——支持从本地 MongoDB 实例将数据批量迁移到远程服务器。它使用 skip + limit 分页策略,每批处理 1000 条记录,并在批次间加入 0.1 秒的延迟,避免对目标数据库造成过大压力:
def migrate_collection(local_collection, remote_collection, batch_size=1000):
total_docs = local_collection.count_documents({})
processed = 0
while processed < total_docs:
batch = list(local_collection.find().skip(processed).limit(batch_size))
result = remote_collection.insert_many(batch)
processed += len(result.inserted_ids)
time.sleep(0.1) # 绅士延迟
这种"遍历 + 分批 + 限速"的模式,在很多数据同步场景中都可以直接复用。
5.4 批量写入优化
controller 中的地址统计场景使用了 bulk_write + upsert,这是 MongoDB 写入性能最高的模式之一。当采集到大量区块时,地址统计会被累积,然后一次性写入,减少网络往返:
bulk_ops = [UpdateOne({'address': addr}, {'$inc': {'transaction_count': count}}, upsert=True)
for addr, count in addresses.items()]
result = collection.bulk_write(bulk_ops)
六、工具模块:Pool 池化设计
项目中还有一个容易被忽略但设计精巧的模块——Pool.py。它提供了四种线程安全的池化数据结构,每个都有明确的适用场景。
6.1 ProxyPool:代理 IP 池
class ProxyPool:
def __init__(self, proxies: List[Dict]):
self._proxies = proxies
self._lock = threading.Lock()
self._cur = 0
def get_proxy(self) -> Optional[Dict]:
with self._lock:
if self._cur == len(self._proxies):
self._cur = 0 # 循环复用
proxy = self._proxies[self._cur]
self._cur += 1
return proxy
这是一个典型的轮询(Round-Robin)代理分配器。与 configuration.py 中的 API Key 随机选择不同,ProxyPool 使用顺序轮询 + 循环复用的策略,确保每个代理被均匀使用。del_proxy 方法在删除代理时还会智能调整内部指针,避免越界或跳过一个代理。
6.2 NumberPool:区块号分发器
class NumberPool:
def __init__(self, start_number: int, end_number: int):
self._cur_numbers = start_number
self._cache_num_list = []
def get_number(self) -> Optional[int]:
with self._lock:
if self._cache_num_list:
return self._cache_num_list.pop()
number = self._cur_numbers
self._cur_numbers += 1
with open("./cache.json", "w") as f:
json.dump({'cur': number}, f)
return number
NumberPool 的设计很有意思——它不仅是一个线程安全的递增计数器,还实现了断点持久化。每次取号都会将当前进度写入 cache.json,这意味着即使程序崩溃重启,也能从上次的位置继续,而不是从头开始。_cache_num_list 则充当了一个"回退缓冲池",当某个区块处理失败时可以放回队列重新调度。
6.3 TokenPool:FIFO Token 管理器
class TokenPool:
def __init__(self, tokens: List[str]):
self._token_queue = deque(tokens)
def get_token(self):
return self._token_queue.popleft() # 从队首取
def put_token(self, token: str):
self._token_queue.append(token) # 放回队尾
TokenPool 使用 collections.deque 实现了一个严格的 FIFO 队列。当某个请求消耗了一个配额 Token 后,用完后放回队尾,实现公平调度。这种模式非常适合第三方 API 的配额管理——确保每个 Token 的使用次数大致均衡。
6.4 Counter:QPS 统计器
class Counter:
def __init__(self):
self._counts = deque()
self._current_count = 0
self._last_second = int(time.time())
def increment(self):
current_second = int(time.time())
if current_second != self._last_second:
self._counts.append(self._current_count)
self._current_count = 1
self._last_second = current_second
else:
self._current_count += 1
def get_speed(self):
counts = list(self._counts)
return sum(counts) / len(counts) if counts else 0
Counter 按秒统计请求次数,可以精确计算 QPS(每秒查询数)。这对于控制请求速率、避免触发 TRON Grid API 的频率限制非常有价值。在实际部署中,可以结合 Counter 的实时 QPS 数据动态调整并发数。
七、并发控制与异步设计
7.1 信号量并发控制
在 controller.py 中,批量获取区块数据时使用了 asyncio.Semaphore 限制并发数:
async def get_range_num_data_async(self, o_num: int, latest_num: int):
semaphore = asyncio.Semaphore(10) # 同时最多 10 个请求
tasks = [
self._fetch_block_with_semaphore(block_num, semaphore)
for block_num in range(o_num, latest_num)
]
results = await asyncio.gather(*tasks)
这不仅保护了目标 API 不被过量请求冲垮,也避免了本地文件描述符耗尽。
7.2 事件循环的生命周期管理
项目在处理同步与异步混用的场景时非常细致。_run_async 方法优雅地处理了三种事件循环状态:
def _run_async(self, coro):
try:
if self._loop and not self._loop.is_closed():
loop = self._loop
else:
loop = asyncio.get_running_loop()
except RuntimeError:
self._loop = asyncio.new_event_loop()
asyncio.set_event_loop(self._loop)
loop = self._loop
if loop.is_running():
future = asyncio.run_coroutine_threadsafe(coro, loop)
return future.result()
else:
return loop.run_until_complete(coro)
这种设计保证了无论在哪种上下文中调用(Flask 请求线程、Monitor 后台线程、测试环境),都能正确执行异步代码。
八、扩展思考:从项目中学到什么
8.1 可借鉴的设计模式
- Scrapy + Flask 双引擎:Scrapy 负责高吞吐的数据采集,Flask 负责轻量控制。二者通过线程解耦,互不干扰。
- 增量批处理:不是一次性拉取所有历史数据,而是按 batch 分片,每批完成后自动查询最新高度,实现了"追块"效果。
- 配置热更新:通过 PUT API 修改配置后立即写入
config.json,无需重启服务即可生效。
8.2 可能的改进方向
- 断点续传持久化:当前爬虫的
start_block是固定的,如果爬虫中途崩溃会丢失进度。可以将当前处理高度写入数据库或 Redis 实现断点续传。Pool.py中的NumberPool已经通过cache.json做了类似的事情,可以将这个思路推广到爬虫主流程中。 - 分布式扩展:引入
scrapy-redis可以将任务队列外移到 Redis,实现多节点分布式采集。Scrapy 的调度器、去重过滤器、队列都可以无缝切换为 Redis 后端,水平扩展能力大大增强。 - 消息队列解耦:爬虫产出数据后不直接写入 MongoDB,而是先发到 Kafka / RabbitMQ,由下游消费者处理,进一步解耦。这样数据采集与数据消费可以独立扩缩容,系统整体的鲁棒性也会更高。
- 更完善的监控告警:当前 Monitor 的异常处理逻辑被注释掉了(
try/except块被注释),可以在生产环境中补全这些逻辑,并结合 Prometheus + Grafana 暴露爬虫运行指标,如当前区块高度、采集延迟、错误率等。 - 配置中心化:
config.json本地文件管理在单机场景下够用,但在多节点部署时建议接入 etcd 或 Consul 等配置中心,实现配置的集中管理与动态推送。
8.1 项目结构中的设计哲学
回顾整个项目,我认为最值得学习的不是某个具体技术点,而是作者解决问题的思路:
- 渐进式设计:从
Pool.py的基础工具 →configuration.py的配置管理 →controller.py的异步封装 →app.py/run.py的 HTTP 控制,每一层都在前一层的肩膀上构建,职责分明,可替换性强。 - 防御性编程:数据库连接有重试、API 请求有重试、地址转换有校验、Item 处理有类型判断,每一处边界都做了防御处理。
- 实用主义优先:没有追求纯异步全链路(Flask + Scrapy 同步混用),没有引入过多中间件,而是选择"够用就好"的方案,让项目保持在"看一眼就能跑起来"的复杂度范围内。
九、快速上手体验
如果你想在本地体验这个项目,只需要几步:
# 1. 克隆项目
git clone https://github.com/Lireal-w/tron_address.git
cd tron_address
# 2. 安装依赖
pip install scrapy flask pymongo aiohttp base58
# 3. 配置 MongoDB 连接和 TRON API Key
cp config.example.json config.json
cp key.example.json key.json
# 编辑 config.json 填入 MongoDB 地址
# 编辑 key.json 填入 TRON Grid API Key
# 4. 启动 Flask 控制服务
python app.py
# 访问 http://localhost:5000
之后就可以通过 POST http://localhost:5000/start_crawler 启动爬虫,或用 GET http://localhost:5000/api/config 查看运行配置了。
总结
tron_address 项目虽然体量不大,但它完整地展示了如何用 Python 生态中成熟的工具链构建一个区块链数据采集系统。Scrapy 提供了高效的异步爬取能力,Flask 提供了灵活的 API 控制层,MongoDB 提供了可靠的持久化存储,而 aiohttp 和 asyncio 则补足了高并发场景下的性能需求。
对于正在搭建区块链数据基础设施的团队来说,这个项目的架构思路——采集与控制分离、增量批处理、API Key 轮换、去重写入——都是值得借鉴的实战经验。希望本文的剖析能为你的链上数据之旅带来一些启发。
评论