Python 连接 etcd
用 Python 连接 etcd 主要依赖 python-etcd3 或 etcd3 这两个库,它们都支持 etcd 的 v3 API。下面的示例会以使用更广泛的 python-etcd3 为主。
📦 准备工作:安装库
首先需要安装 python-etcd3 库。
pip install etcd3示例
import os
os.environ["PROTOCOL_BUFFERS_PYTHON_IMPLEMENTATION"] = "python"
import etcd3
def main():
# 连接到本地默认端口
client = etcd3.client(host='localhost', port=2379)
client.put('/service/config', 'hello world')
value, metadata = client.get('/service/config')
print(f'值: {value}, 元数据: {metadata}')
if __name__ == '__main__':
main()输出结果
值: b'hello world', 元数据: <etcd3.client.KVMetadata object at 0x106ef0670>🔌 一、建立连接
1. 连接单个节点
这是最基础的用法,连接本地默认的 etcd 服务。
import etcd3
# 连接到本地默认端口
client = etcd3.client(host='localhost', port=2379)2. 连接集群
在生产环境中,通常会部署一个 etcd 集群。客户端可以连接集群中的任意一个节点。
import etcd3
# 连接集群中的任意一个节点
client = etcd3.client(host='192.168.1.101', port=2379)3. 带认证的连接
如果 etcd 服务启用了安全认证,需要在建立连接时提供用户名和密码。
import etcd3
client = etcd3.client(
host='localhost',
port=2379,
user='myuser', # 用户名
password='mypass' # 密码
)📝 二、基础操作 (CRUD)
1. 写入 (Put)
写入键值对。
# 写入一个简单的键值对
client.put('/service/config', 'hello world')2. 读取 (Get)
读取指定键的值。
# 读取单个键
value, metadata = client.get('/service/config')
print(f'值: {value}, 元数据: {metadata}')
# 输出: 值: b'hello world', 元数据: <KeyMeta object at ...>get 方法返回的是一个包含值和元数据的元组,元数据里包含了版本、创建/修改的 revision 等信息。如果键不存在,则 value 为 None。
3. 批量读取 (Get with Prefix)
读取所有以某个前缀开头的键。
# 读取所有以 '/service/' 为前缀的键
results = client.get_prefix('/service/')
for value, metadata in results:
print(f'键: {metadata.key.decode()}, 值: {value.decode()}')4. 删除 (Delete)
删除一个或多个键。
# 删除单个键
client.delete('/service/config')
# 删除所有以 '/service/' 为前缀的键
client.delete_prefix('/service/')🛠️ 三、高级特性
1. 租约 (Lease)
租约可以为键值对设置一个有效期(TTL)。当租约过期时,所有与之绑定的键值对都会被自动删除。
# 创建一个 TTL 为 60 秒的租约
lease = client.lease(60)
# 将键绑定到该租约
client.put('/temp/data', 'this will disappear in 60 seconds', lease=lease)
# 手动刷新租约,重置 TTL
lease.refresh()
# 获取租约的剩余 TTL
print(lease.ttl) # 输出剩余的秒数2. 监听 (Watch)
监听一个或多个键的变化。当被监听的键被修改、删除时,客户端会收到事件通知。
# 监听单个键
events_iterator, cancel = client.watch('/service/config')
for event in events_iterator:
# 当 '/service/config' 发生变化时,会打印出事件详情
print(event)
# 如果只想监听一次,可以调用 cancel() 停止监听
cancel()
break3. 原子操作 (Transaction)
事务可以保证多个操作要么全部成功,要么全部失败,是保证数据一致性的重要手段。
# 创建一个事务对象
txn = client.transaction()
# 添加条件:如果键 '/my/key' 的值为 'old_value'
txn.compare(value='/my/key') == 'old_value'
# 条件成立时执行的操作:更新键值
txn.success(put='/my/key', value='new_value')
# 条件不成立时执行的操作:什么都不做
txn.failure(put='/my/key', value='unexpected')
# 提交事务
txn.commit()4. 分布式锁 (Lock)
使用 etcd 的分布式锁,可以确保在分布式系统中,同一时间只有一个客户端(或服务)能执行某段关键代码。
# 创建一个名为 'my-resource-lock' 的锁
lock = client.lock('my-resource-lock')
# 获取锁
if lock.acquire(timeout=10):
try:
# 成功获取锁,执行需要互斥的业务逻辑
print("Lock acquired, doing critical work...")
# ... 你的代码 ...
finally:
# 最后一定要释放锁
lock.release()
else:
print("Failed to acquire lock within timeout.")
# 也可以使用上下文管理器(with 语句)来自动管理锁
with client.lock('another-resource-lock', ttl=30) as lock:
# 在此代码块中,锁已被成功获取
print("Working inside the lock...")
# 代码块结束时,锁会自动释放🎯 四、实用示例:服务注册与心跳
下面是一个完整的例子,展示了如何用 etcd 实现服务注册与心跳。一个服务启动时,会在 etcd 中注册一个临时节点并持续发送心跳来续约,确保 etcd 能准确反映服务的在线状态。
import etcd3
import time
import threading
class ServiceRegistry:
def __init__(self, etcd_host, etcd_port, service_name, service_addr, ttl=10):
self.client = etcd3.client(host=etcd_host, port=etcd_port)
self.service_name = service_name
self.service_addr = service_addr
self.ttl = ttl
self.lease = None
self.keep_alive_thread = None
self.running = False
def register(self):
"""注册服务"""
self.lease = self.client.lease(self.ttl)
key = f'/services/{self.service_name}/{self.service_addr}'
self.client.put(key, 'online', lease=self.lease)
print(f'Service {self.service_name} registered at {self.service_addr}')
self.start_keep_alive()
def start_keep_alive(self):
"""启动心跳线程,定期刷新租约"""
self.running = True
self.keep_alive_thread = threading.Thread(target=self._keep_alive_loop)
self.keep_alive_thread.daemon = True
self.keep_alive_thread.start()
def _keep_alive_loop(self):
while self.running:
time.sleep(self.ttl / 2)
if self.lease:
try:
# 刷新租约,重置TTL
self.lease.refresh()
except Exception as e:
print(f"Failed to refresh lease: {e}")
def unregister(self):
"""注销服务"""
self.running = False
key = f'/services/{self.service_name}/{self.service_addr}'
self.client.delete(key)
if self.keep_alive_thread:
self.keep_alive_thread.join()
print(f'Service {self.service_name} unregistered')
# 使用示例
if __name__ == '__main__':
registry = ServiceRegistry('localhost', 2379, 'web-api', 'http://10.0.0.1:8080')
try:
registry.register()
print("Service is running, press Ctrl+C to stop...")
# 让主线程持续运行
while True:
time.sleep(1)
except KeyboardInterrupt:
registry.unregister()⚠️ 五、重要注意事项
- API 版本:确保你使用的是 v3 的 API 客户端。像
python-etcd(对应 etcd v2)已经是过时的版本,不应再使用。 - 键的命名规范:建议使用
string类型作为键,并采用如/service/name这样的层级结构,便于管理和查询。 - 连接超时:在生产环境中,应考虑设置连接超时,避免因网络问题导致程序无限等待。
- 上下文管理器:对于
Lock这类资源,尽量使用with语句来管理,确保资源在任何情况下都能被正确释放。 - 错误处理:实际使用时,务必为网络请求(如
put、get、lease.refresh())添加try...except块,以处理可能出现的连接错误或超时等异常。
在理解了这些基本用法和重要注意事项之后,如果你有兴趣尝试 etcd3-py 或 etcd3gw 等替代库,或者想深入探讨某个具体的实现细节,随时都可以提出来~
报错
问题描述
TypeError: Descriptors cannot be created directly.
If this call came from a _pb2.py file, your generated code is out of date and must be regenerated with protoc >= 3.19.0.
If you cannot immediately regenerate your protos, some other possible workarounds are:
1. Downgrade the protobuf package to 3.20.x or lower.
2. Set PROTOCOL_BUFFERS_PYTHON_IMPLEMENTATION=python (but this will use pure-Python parsing and will be much slower).
More information: https://developers.google.com/protocol-buffers/docs/news/2022-05-06#python-updates这个错误是因为你的 protobuf Python 包版本过新(≥4.21.0),而项目中的 _pb2.py 文件是由旧版 protoc 生成的,两者不兼容。
解决方案(三选一)
1. ✅ 推荐:重新生成 _pb2.py 文件(保持新版 protobuf)
- 升级
protoc到 ≥3.19.0
# 检查当前版本
protoc --version
# 如果版本过低,从 https://github.com/protocolbuffers/protobuf/releases 下载安装- 重新编译
.proto文件
protoc --python_out=. your_file.proto- 重新生成后,错误即可消失,且能使用 protobuf 最新特性。
2. 降级 protobuf 包到 3.20.x
pip install protobuf==3.20.3⚠️ 这会锁定旧版本,可能影响其他依赖。
3. 设置环境变量(纯 Python 解析,性能较低)
export PROTOCOL_BUFFERS_PYTHON_IMPLEMENTATION=python # Linux/macOS
set PROTOCOL_BUFFERS_PYTHON_IMPLEMENTATION=python # Windows或在代码开头添加:
import os
os.environ["PROTOCOL_BUFFERS_PYTHON_IMPLEMENTATION"] = "python"为什么会出现这个错误?
- protobuf 4.x 改变了描述符的创建方式,旧版生成的
_pb2.py中使用了被移除的内部 API。 - 如果你无法重新生成 proto 文件(例如来自第三方库),降级是更快的临时方案。
建议优先尝试方案 1,能彻底解决问题并保持环境最新。