diff --git "a/Django/\345\237\272\347\241\200\346\225\231\347\250\213/\345\205\245\351\227\250\347\233\270\345\205\263/\346\212\225\347\245\250\345\272\224\347\224\250-\346\250\241\345\236\213.md" "b/Django/\345\237\272\347\241\200\346\225\231\347\250\213/\345\205\245\351\227\250\347\233\270\345\205\263/\346\212\225\347\245\250\345\272\224\347\224\250-\346\250\241\345\236\213.md" index 3a87c10..fd42fe8 100644 --- "a/Django/\345\237\272\347\241\200\346\225\231\347\250\213/\345\205\245\351\227\250\347\233\270\345\205\263/\346\212\225\347\245\250\345\272\224\347\224\250-\346\250\241\345\236\213.md" +++ "b/Django/\345\237\272\347\241\200\346\225\231\347\250\213/\345\205\245\351\227\250\347\233\270\345\205\263/\346\212\225\347\245\250\345\272\224\347\224\250-\346\250\241\345\236\213.md" @@ -16,7 +16,8 @@ mysite/ # 项目目录 # 应用集配置 ```bash -vim mysite/mysite/settings.py +cd mysite +vim mysite/settings.py CUSTOMIZED_APPS = [ 'mysite', # 自定义的应用 @@ -41,7 +42,8 @@ INSTALLED_APPS = INSTALLED_APPS + [ # 数据库配置 ```bash -vim mysite/mysite/settings.py +cd mysite +vim mysite/settings.py DATABASES = { 'default': { # 默认使用的数据库,可配置多个选用 @@ -73,6 +75,7 @@ python manage.py runserver 0.0.0.0:8000 ```bash python manage.py startapp polls +# 原结构 tree polls/ polls/ # 自定义应用名 ├── __init__.py # 表示应用是一个包 @@ -83,8 +86,139 @@ polls/ # 自定义应用名 ├── models.py # 数据库驱动应用的模型文件,深度定制通常创建models包,内部再对每个模型独立定制化 ├── tests.py # 单元测试文件,深度定制通常创建tests包,内部再对每个模型独立定制化 └── views.py # 视图处理文件,深度定制通常创建views包,内部再对每个模型独立定制化 + +# 现结构 + ``` * 应用是一个Web应用程序(必须是一个Python包),它完成具体事项,如博客应用,投票应用等,项目是相关配置和应用的集合,一个项目可以包含多个应用,一个应用也可以被打包分发给不同的项目使用. * 通过startapp命令可以自动生成应用的基本目录结构,由于manage.py与应用同级,所以可以在代码中任意位置直接引用这些应用甚至自定义的Python包. +# 创建应用模型 + +```python +cd mysite +vim polls/models.py + +class Question(models.Model): + question_text = models.CharField(max_length=200) + pub_date = models.DateTimeField('date published') + + +class Choice(models.Model): + question = models.ForeignKey(Question) + choice_text = models.CharField(max_length=200) + votes = models.IntegerField(default=0) +``` + +* 编写数据库驱动的应用第一步就是定义模型,如上创建2个模型,问题模型Question和选项模型Choice,前者包含一个问题字段和发布时间字段,后者包含一个选项内容字段和得票字段,每个Choice实例都与Question实例外键关联 +* 模型类都继承自django.db.models.Model,对应数据库中的表,类变量即模型字段都继承自django.db.models.fields.Field,对应数据库中表字段,ForeignKey对应数据库中表外建字段. + +# 激活应用模型 + +```bash +cd mysite +vim mysite/settings.py + +CUSTOMIZED_APPS = [ + 'mysite', + 'polls', +] + +INSTALLED_APPS = CUSTOMIZED_APPS + [ + 'django.contrib.admin', + 'django.contrib.auth', + 'django.contrib.contenttypes', + 'django.contrib.sessions', + 'django.contrib.messages', + 'django.contrib.staticfiles', +] + +INSTALLED_APPS = INSTALLED_APPS + [ + +] +``` + +* 注意只有添加到入口应用配置中的INSTALLED_APPS中的应用才会被激活,也就是应用下面的模型才会被接管,所以Django应用支持“热插拔”,可以分发给其它项目使用或多个项目本地公用 + +```bash +cd mysite +python manage.py makemigrations polls + +Migrations for 'polls': + polls/migrations/0001_initial.py + - Create model Choice + - Create model Question + - Add field question to choice +``` + +* 通过运行makemigrations 告知Django应用polls的模型发生了改变,需要生成迁移脚本 + +```bash +cd mysite +python manage.py sqlmigrate polls 0001 + +BEGIN; +-- +-- Create model Choice +-- +CREATE TABLE "polls_choice" ("id" integer NOT NULL PRIMARY KEY AUTOINCREMENT, "choice_text" varchar(200) NOT NULL, "votes" integer NOT NULL); +-- +-- Create model Question +-- +CREATE TABLE "polls_question" ("id" integer NOT NULL PRIMARY KEY AUTOINCREMENT, "question_text" varchar(200) NOT NULL, "pub_date" datetime NOT NULL); +-- +-- Add field question to choice +-- +ALTER TABLE "polls_choice" RENAME TO "polls_choice__old"; +CREATE TABLE "polls_choice" ("id" integer NOT NULL PRIMARY KEY AUTOINCREMENT, "choice_text" varchar(200) NOT NULL, "votes" integer NOT NULL, "question_id" integer NOT NULL REFERENCES "polls_question" ("id")); +INSERT INTO "polls_choice" ("choice_text", "votes", "id", "question_id") SELECT "choice_text", "votes", "id", NULL FROM "polls_choice__old"; +DROP TABLE "polls_choice__old"; +CREATE INDEX "polls_choice_question_id_c5b4b260" ON "polls_choice" ("question_id"); +COMMIT; +``` + +* 通过运行sqlmigrate 获取迁移脚本转换后的SQL语句 + +```bash +cd mysite +python manage.py check +System check identified no issues (0 silenced). +``` + +* 通过运行check命令检查激活的应用模型是否存在问题而并不会执行迁移脚本 + +```bash +python manage.py migrate + +Operations to perform: + Apply all migrations: admin, auth, contenttypes, polls, sessions +Running migrations: + Applying polls.0001_initial... OK +``` + +* 通过运行migrate会对比数据库中迁移版本号将应用下migrations目录下需要迁移的脚本转换为SQL执行 + +* 迁移功能允许你在开发过程中不断修改模型而不用删除数据库或表,[更多功能](#http://www.2xkt.com/documents/django_182/ref/django-admin.html) + +# 模型访问接口 + +````bash +cd mysite +python manage.py shell +>>> +或 +cd mysite +python +>>> import os +>>> os.environ.setdefault("DJANGO_SETTINGS_MODULE", "mysite.settings") +>>> import django +>>> django.setup() +```` + +* 通过运行shell会进入Django API交互环境,它会加载入口应用下的应用集配置settings.py,当然也可以在默认Python shell中通过设置环境变量DJANGO_SETTINGS_MODULE再调用django.setup()进入 + +```python + +``` + diff --git "a/Django/\346\241\206\346\236\266\351\242\206\346\202\237/\350\256\276\350\256\241\345\223\262\345\255\246.md" "b/Django/\346\241\206\346\236\266\351\242\206\346\202\237/\350\256\276\350\256\241\345\223\262\345\255\246.md" new file mode 100644 index 0000000..e785589 --- /dev/null +++ "b/Django/\346\241\206\346\236\266\351\242\206\346\202\237/\350\256\276\350\256\241\345\223\262\345\255\246.md" @@ -0,0 +1,2 @@ +http://www.2xkt.com/documents/django_182/misc/design-philosophies.html#dry + diff --git "a/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\345\257\271\350\261\241\346\211\251\345\261\225/contextlib.md" "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\345\257\271\350\261\241\346\211\251\345\261\225/contextlib.md" new file mode 100644 index 0000000..a8f4a43 --- /dev/null +++ "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\345\257\271\350\261\241\346\211\251\345\261\225/contextlib.md" @@ -0,0 +1,114 @@ +---- + +* [模块由来](#模块由来) +* [模块简介](#模块简介) +* [实现方式](#实现方式) + * [基于类的实现](#基于类的实现) + * [基于生成器+contextlib.contextmanager装饰器实现](#基于生成器+contextlib.contextmanager装饰器实现) +* [属性方法](#属性方法) + +---- + +# 模块由来 + +> 通常希望在某语句块儿运行直至结束保持某种状态(上下文),精确的分配和释放资源. + +# 模块简介 + +> 内置模块,简化上下文管理器的声明来支持with语句 + +# 实现方式 + +> 上下文管理器的实现方式通常可以基于类的\_\_exit\_\_和\_\_enter\_\_实现,也可以基于yield生成器和contextlib.contextmanager装饰器实现 + +## 基于类的实现 + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +class FileWithContextManager(object): + def __init__(self, *args, **kwargs): + self.fobj = file(*args, **kwargs) + + def __enter__(self): + return self.fobj + + def __exit__(self, exc_type, exc_val, exc_tb): + if not self.fobj.closed: self.fobj.close() + + +if __name__ == '__main__': + with FileWithContextManager('/etc/passwd') as f: + for line in f: + print line, +``` + +* with语句首先暂存了FileWithContextManager的\_\_exit\_\_方法,然后调用\_\_enter\_\_返回给with通过as赋值给变量f,当with内部语句块儿执行完毕时调用之前暂存的\_\_exit\_\_方法 +* 如果语句块儿内异常则会将exc_type, exc_val, exc_tb传递给\_\_exit\_\_方法,但需要注意的是_\_exit\_\_如果返回True则会优雅忽略,否则会被重新抛出 + +## 基于生成器+contextlib.contextmanager装饰器实现 + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +from contextlib import contextmanager + + +@contextmanager +def open_file(*args, **kwargs): + try: + f = open(*args, **kwargs) + yield f + except Exception as e: + f = None + raise e + finally: + if f and not f.closed: f.close() + + +if __name__ == '__main__': + with open_file('/etc/passwd') as f: + for line in f: + print line, +``` + +* contextmanager其实内部封装了open\_file(*args, **kwargs)为GeneratorContextManager对象,此对象在其\_\_enter\_\_中调用open\_file(*args, **kwargs).next()方法,在\_\_exit\_\_中再次调用open_file(*args, **kwargs).next()方法,所以被修饰的函数或方法必须是一个生成器 + +# 属性方法 + +| 方法 | 说明 | +| -------------------- | ------------------------------------------------------------ | +| closing(thing) | 配合with..as语句在语句块儿执行完毕后调用被封装对象的.close() | +| nested(*managers) | 配合with..as语句同时创建多个上下文管理器 | +| contextmanager(func) | 作为装饰器装饰被yield分隔为2部分的函数或方法或对象,yield前交由\_\_enter\_\_处理,yield后交由\_\_exit\_\_处理,yield的值可通过as传递 | + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import contextlib + + +if __name__ == '__main__': + # < Python2.7 + with contextlib.nested(open('/etc/passwd', 'rb'), open('passwd', 'ab')) as (src_passwd, dst_passwd): + for line in src_passwd: + dst_passwd.write(line) + + # >= Python2.7 + with open('/etc/passwd', 'rb') as src_passwd, open('passwd', 'ab') as dst_passwd: + for line in src_passwd: + dst_passwd.write(line) +``` + diff --git "a/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\345\257\271\350\261\241\346\211\251\345\261\225/itertools.md" "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\345\257\271\350\261\241\346\211\251\345\261\225/itertools.md" new file mode 100644 index 0000000..f855018 --- /dev/null +++ "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\345\257\271\350\261\241\346\211\251\345\261\225/itertools.md" @@ -0,0 +1,13 @@ +# 迭代器有意思的例子 + +```python +@staticmethod +def _get_tasks(func, it, size): + it = iter(it) + while 1: + x = tuple(itertools.islice(it, size)) + if not x: + return + yield (func, x) +``` + diff --git "a/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\346\225\260\346\215\256\346\240\241\351\252\214/crcmod.md" "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\346\225\260\346\215\256\346\240\241\351\252\214/crcmod.md" index bab953c..f6e6095 100644 --- "a/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\346\225\260\346\215\256\346\240\241\351\252\214/crcmod.md" +++ "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\346\225\260\346\215\256\346\240\241\351\252\214/crcmod.md" @@ -1,7 +1,3 @@ -# 模块简介 - - - # 循环冗余校验 ## 基本原理 @@ -54,5 +50,3 @@ CRC校验的前提是客户端和服务端事先约定一个模二除法中的除数(生成多项式),除数的最高位与最低位必须是1,每一个多项式都与一个二进制除数对应,如CRC-8表示100110001,生成多项式主要用于降低校验出错率 - - diff --git "a/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\346\250\241\345\235\227\345\257\274\345\205\245/imp.md" "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\346\250\241\345\235\227\345\257\274\345\205\245/imp.md" new file mode 100644 index 0000000..42204d0 --- /dev/null +++ "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\346\250\241\345\235\227\345\257\274\345\205\245/imp.md" @@ -0,0 +1,57 @@ + + +---- + +* [模块简介](#模块简介) +* [属性方法](#属性方法) +* [实战练习](#实战练习) + +---- + +# 模块简介 + +> 内置模块,默认提供import语句实现,内部依赖importlib,从Python3.4已废弃,相关方法已迁移至importlib + +# 属性方法 + +| 属性 | 说明 | +| ------------------------- | ------------------------------------------------------------ | +| find_module(name, [path]) | path未指定则从sys.path查找name,否则从path查找name,path通常为module_obj.\_\_path\_\_,如果未找到抛出ImportError异常 | + + + +# 实战练习 + +* 查询django.contrib.sessions包下是否包含middleware模块 ? + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import imp +import importlib + + +def has_module(package, name): + # 尝试导入导入字符串表示的包 + try: + pkg_path = importlib.import_module(package).__path__ + except AttributeError: + return False + + # 在模块所在的路径下查找指定模块 + try: + imp.find_module(name, pkg_path) + except ImportError: + return False + + return True + + +if __name__ == '__main__': + print has_module('django.contrib.sessions', 'middleware') +``` + diff --git "a/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\346\250\241\345\235\227\345\257\274\345\205\245/importlib.md" "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\346\250\241\345\235\227\345\257\274\345\205\245/importlib.md" new file mode 100644 index 0000000..533784a --- /dev/null +++ "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\346\250\241\345\235\227\345\257\274\345\205\245/importlib.md" @@ -0,0 +1,245 @@ +---- + +* [模块简介](#模块简介) +* [属性方法](#属性方法) +* [应用场景](#应用场景) + * [声明位置](#声明位置) + * [调用位置](#调用位置) +* [实战练习](#实战练习) +* [实战总结](#实战总结) + +---- + +# 模块简介 + +> 内置模块,默认提供import语句和\_\_import\_\_底层实现,包括动态导入,导入检查等特性 + +# 属性方法 + +| 属性 | 说明 | +| :-------------------------------- | ------------------------------------------------------------ | +| import_module(name, package=None) | dotted_path导入模块,当name为.开头的相对导入字符串,package必须存在 | + +# 应用场景 + +> Django框架中自动加载应用集配置中的中间件. + +## 声明位置 + +> django.utils.module_loading + +```python +# 用于导入点连接的导入字符串,例如django.contrib.sessions.middleware.SessionMiddleware +def import_string(dotted_path): + """ + Import a dotted module path and return the attribute/class designated by the + last name in the path. Raise ImportError if the import failed. + """ + try: + # 首先按.分隔,如上 + # module_path为django.contrib.sessions.middleware + # class_name为SessionMiddleware + module_path, class_name = dotted_path.rsplit('.', 1) + except ValueError: + msg = "%s doesn't look like a module path" % dotted_path + # 主要为了兼容Py2和Py3,其实就是调用raise重新抛出了ImportError + six.reraise(ImportError, ImportError(msg), sys.exc_info()[2]) + # 调用import_module直接导入点连接module_path返回模块对象 + module = import_module(module_path) + + try: + # 尝试运行时反省获取模块中名称为class_name的类 + return getattr(module, class_name) + except AttributeError: + # 重新抛出异常,依然定义为导入错误ImportError + msg = 'Module "%s" does not define a "%s" attribute/class' % ( + module_path, class_name) + six.reraise(ImportError, ImportError(msg), sys.exc_info()[2]) +``` + +## 调用位置 + +> django.core.handlers.base + +```python +class BaseHandler(object): + + def __init__(self): + self._request_middleware = None + self._view_middleware = None + self._template_response_middleware = None + self._response_middleware = None + self._exception_middleware = None + self._middleware_chain = None + # 加载配置中的中间件 + def load_middleware(self): + """ + Populate middleware lists from settings.MIDDLEWARE (or the deprecated + MIDDLEWARE_CLASSES). + + Must be called after the environment is fixed (see __call__ in subclasses). + """ + self._request_middleware = [] + self._view_middleware = [] + self._template_response_middleware = [] + self._response_middleware = [] + self._exception_middleware = [] + + if settings.MIDDLEWARE is None: + warnings.warn( + "Old-style middleware using settings.MIDDLEWARE_CLASSES is " + "deprecated. Update your middleware and use settings.MIDDLEWARE " + "instead.", RemovedInDjango20Warning + ) + handler = convert_exception_to_response(self._legacy_get_response) + for middleware_path in settings.MIDDLEWARE_CLASSES: + # 尝试导入settings.MIDDLEWARE_CLASSES中定义的中间件 + mw_class = import_string(middleware_path) + try: + mw_instance = mw_class() + except MiddlewareNotUsed as exc: + if settings.DEBUG: + if six.text_type(exc): + logger.debug('MiddlewareNotUsed(%r): %s', middleware_path, exc) + else: + logger.debug('MiddlewareNotUsed: %r', middleware_path) + continue + + if hasattr(mw_instance, 'process_request'): + self._request_middleware.append(mw_instance.process_request) + if hasattr(mw_instance, 'process_view'): + self._view_middleware.append(mw_instance.process_view) + if hasattr(mw_instance, 'process_template_response'): + self._template_response_middleware.insert(0, mw_instance.process_template_response) + if hasattr(mw_instance, 'process_response'): + self._response_middleware.insert(0, mw_instance.process_response) + if hasattr(mw_instance, 'process_exception'): + self._exception_middleware.insert(0, mw_instance.process_exception) + else: + handler = convert_exception_to_response(self._get_response) + for middleware_path in reversed(settings.MIDDLEWARE): + # 尝试导入settings.MIDDLEWARE中定义的中间件 + middleware = import_string(middleware_path) + try: + mw_instance = middleware(handler) + except MiddlewareNotUsed as exc: + if settings.DEBUG: + if six.text_type(exc): + logger.debug('MiddlewareNotUsed(%r): %s', middleware_path, exc) + else: + logger.debug('MiddlewareNotUsed: %r', middleware_path) + continue + + if mw_instance is None: + raise ImproperlyConfigured( + 'Middleware factory %s returned None.' % middleware_path + ) + + if hasattr(mw_instance, 'process_view'): + self._view_middleware.insert(0, mw_instance.process_view) + if hasattr(mw_instance, 'process_template_response'): + self._template_response_middleware.append(mw_instance.process_template_response) + if hasattr(mw_instance, 'process_exception'): + self._exception_middleware.append(mw_instance.process_exception) + + handler = convert_exception_to_response(mw_instance) + + # We only assign to this when initialization is complete as it is used + # as a flag for initialization being complete. + self._middleware_chain = handler +``` + + + +# 实战练习 + +* 如何实现项目中模块和包的递归自动导入 ? + * 思考 + * 怎么配合元类实现自动注册功能 ? + * 尝试模块化Django应用的后台,模型,视图,测试并实现自动导入 ? + +> mysite/utils/module_loading.py + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import os +import imp +import importlib + + +def import_sub_module(package, name): + # 第三方包默认都存在__path__,如果不存在则说明传参错误 + try: + m = importlib.import_module(package) + path = m.__path__ + except AttributeError: + return + # 从包路径查找模块或包 + try: + imp.find_module(name, path) + except ImportError: + return + # 按照绝对路径导入字符串导入 + dotted_path = '{0}.{1}'.format(package, name) + return importlib.import_module(dotted_path) + + +def autodiscovery_modules(package, entrance): + # 记录导入的模块对象 + modules = [] + + cur_dir = os.path.dirname(entrance) + pyfiles = os.listdir(cur_dir) + # 遍历当前目录尝试导入目录下模块和包 + for f_name in pyfiles: + f_path = os.path.join(cur_dir, f_name) + if os.path.isfile(f_path): + if not f_name.endswith('.py'): + continue + m_name, _, _ = f_name.rpartition('.') + if m_name == '__init__': + continue + else: + # 假设目录就是Python包 + m_name = f_name + # 尝试导入 + m = import_sub_module(package, m_name) + if not m: + continue + modules.append(m) + return modules +``` + +> mysite/polls/models/\_\_init\_\_.py + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +from functools import partial +from utils.module_loading import autodiscovery_modules + +# 当尝试导入此包下的模块或包时就会自动预导入,由于Django默认模型使用元类注册,所以导入时就会自动到基类注册,所以makemigrations不受影响 +modules = autodiscovery_modules(__name__, __file__) + + +# 注入全局变量,并不推荐,只是模拟from x import *,让Django模型和原来一样类似from polls.models import Question, Choice一样正常调用 +g_data = {} +map(lambda m: g_data.update(m.__dict__), modules) +globals().update(g_data) + +autodiscovery = partial(autodiscovery_modules,__name__, __file__) +``` + +# 实战总结 + +* 自动导入通常配合元类注册才能发挥其优势,如上由于Django框架的模型基于元类注册所以才能完美配合自动导入使用 + diff --git "a/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\347\275\221\347\273\234\351\200\232\344\277\241/socket.md" "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\347\275\221\347\273\234\351\200\232\344\277\241/socket.md" new file mode 100644 index 0000000..3ea6632 --- /dev/null +++ "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\347\275\221\347\273\234\351\200\232\344\277\241/socket.md" @@ -0,0 +1,174 @@ +# 常规套路 + +## 服务端 + +```python +#! -*- coding: utf-8 -*- + + +import crcmod +import socket +import struct +import logging + + +logging.basicConfig(level=logging.DEBUG, + format='%(asctime)s - %(name)s - %(pathname)s - %(lineno)d - %(levelname)s - %(message)s') +logger = logging.getLogger(__name__) + + +class BaseServer(object): + def __init__(self, **kwargs): + self.sock = None + self.host = kwargs.get('host', '0.0.0.0') + self.port = kwargs.get('port', 51314) + self.addr = (self.host, self.port) + self.handler = kwargs.get('handler', self._handler) + + def _handler(self, data): + logger.debug('Recv data {0}'.format(data)) + + @staticmethod + def _sock_initopts(sock, opts=None): + user_opts = set() + user_opts.update([ + (socket.SOL_SOCKET, socket.SO_REUSEADDR, 1), + (socket.SOL_SOCKET, socket.SO_REUSEPORT, 1), + ]) + if isinstance(opts, (set, list, tuple)): + user_opts.update(opts) + map(lambda opt: sock.setsockopt(*opt), user_opts) + + def run(self): + raise NotImplementedError + + +class TcpServer(BaseServer): + MODE = 'TcpServer' + + def __init__(self, **kwargs): + super(TcpServer, self).__init__(**kwargs) + self.listen = kwargs.get('listen', 1024) + self.buffer = kwargs.get('buffer', 1024) + + self.sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + self._sock_initopts(self.sock) + self.sock.bind(self.addr) + self.sock.listen(self.listen) + + def run(self): + pass + + +class CrcTcpServer(TcpServer): + MODE = 'CrcTcpServer' + PACK_HEAD_SIZE = 12 + + def __init__(self, **kwargs): + super(CrcTcpServer, self).__init__(**kwargs) + + def run(self): + data_buffer = bytes() + while True: + logger.debug('Wait for connection...') + conn, addr = self.sock.accept() + logger.debug('Connection from {0}'.format(addr)) + + while True: + data = conn.recv(self.buffer) + if not data: + break + data_buffer += data + while True: + data_buffer_size = len(data_buffer) + if data_buffer_size < self.PACK_HEAD_SIZE: + logger.debug('Data buffer size {0} < {1}, break'.format(data_buffer_size, + self.PACK_HEAD_SIZE)) + break + version, bodysize, crcsize = struct.unpack('!3I', data_buffer[:self.PACK_HEAD_SIZE]) + total_size = self.PACK_HEAD_SIZE + bodysize + if data_buffer_size < total_size: + logger.debug('Data buffer size {0} < {1}, break'.format(data_buffer_size, + total_size)) + break + checkdata = data_buffer[self.PACK_HEAD_SIZE:total_size] + print '=' * 100 + print bin(int(checkdata.encode('hex'), 16)) + print '=' * 100 + real_data = data_buffer[self.PACK_HEAD_SIZE:total_size-crcsize] + self.handler(real_data) + data_buffer = data_buffer[total_size:] + + +class Reciver(object): + def __init__(self, **kwargs): + self.server_classes = {} + + def create_server(self, mode, context=None): + context = context or {} + assert mode in self.server_classes, '{0} server klass not registed'.format(mode) + server_class = self.server_classes[mode] + return server_class(**context) + + def register_server_classes(self, *klasses): + for klass in klasses: + if not hasattr(klass, 'MODE'): + print 'MODE undefined in {0}, ignore'.format(klass.__name__) + continue + self.server_classes[klass.MODE] = klass + +if __name__ == '__main__': + crc_tcp_server_context = { + 'host': '0.0.0.0', + 'port': 51314, + } + reciver = Reciver() + reciver.register_server_classes(CrcTcpServer) + server = reciver.create_server('CrcTcpServer', crc_tcp_server_context) + server.run() +``` + +## 客户端 + +```python +#! -*- coding: utf-8 -*- + + +import json +import crcmod +import socket +import struct +import logging +import binascii + + +logging.basicConfig(level=logging.DEBUG, + format='%(asctime)s - %(name)s - %(pathname)s - %(lineno)d - %(levelname)s - %(message)s') +logger = logging.getLogger(__name__) + +if __name__ == '__main__': + sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + sock.connect(('127.0.0.1', 51314)) + + # real_data = json.dumps({'name': u'李满满', 'age': 27}) + + real_data = 'a' + crc32 = crcmod.predefined.Crc('crc-8') + hex_real_data = binascii.hexlify(real_data) + crc32.update(hex_real_data) + real_data_crc = crc32.hexdigest() + real_data_crc_bytes = real_data_crc.decode('hex') + body_with_crc_hex = hex_real_data+real_data_crc + body_with_crc_bytes = body_with_crc_hex.decode('hex') + + version, bodysize, crcsize = 1, len(body_with_crc_bytes), len(real_data_crc_bytes) + head_packed = struct.pack('!3I', *[version, bodysize, crcsize]) + + data = head_packed+body_with_crc_bytes + sock.sendall(data) +``` + +# 实战练习 + +* [编写SOCKET5服务器](#https://hatboy.github.io/2018/04/28/Python%E7%BC%96%E5%86%99socks5%E6%9C%8D%E5%8A%A1%E5%99%A8/) + diff --git "a/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\231\232\346\213\237\347\216\257\345\242\203/pipenv.md" "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\231\232\346\213\237\347\216\257\345\242\203/pipenv.md" new file mode 100644 index 0000000..f142047 --- /dev/null +++ "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\231\232\346\213\237\347\216\257\345\242\203/pipenv.md" @@ -0,0 +1,101 @@ +---- + +* [模块简介](#模块简介) +* [安装部署](#安装部署) +* [常规套路](#常规套路) + * [全局环境变量设置](#全局环境变量设置) + * [创建独立虚拟环境](#创建独立虚拟环境) + * [管理第三方依赖包](#管理第三方依赖包) + * [更新PYPI源的地址](#更新PYPI源的地址) + * [利用虚拟环境开发](#利用虚拟环境开发) +* [错误汇总](#错误汇总) + +---- + +# 模块简介 + +> 第三方模块,可实现一台主机上轻松安装管理切换Python虚拟开发环境 + +* 此模块是Python官方推荐的包管理工具,综合了virtualenv,pip,pyenv三者的功能,使用Pipfile和Pipfile.lock来自动管理依赖包,当通过pipenv添加或删除包时,它会自动维护Pipfile文件并同时生成Pipfile.lock来锁定安装包的版本和依赖信息,避免构建错误 + +# 安装部署 + +```bash +pip install pipenv +``` + +# 常规套路 + +## 全局环境变量设置 + +```bash +echo -e '# pipenv\nexport PIPENV_VENV_IN_PROJECT=1' >> ~/.bash_profile +source ~/.bash_profile +``` + +* 设置在每个项目根目录下创建虚拟环境.venv + +## 创建独立虚拟环境 + +```python +mkdir myproject +pipenv --python ~/.pyenv/versions/3.6.6/bin/python +``` + +## 管理第三方依赖包 + +> pipenv操作指令和pip完全一致,不再重复说明 + +```bash +pipenv install django==1.11.5 +# 安装MySQL-python For Mac, 其它系统比较简单,略 +# ---可能出现的问题 +# fatal error: 'my_config.h' file not found +# --- +# 已安装可跳过 +brew install mysql +brew unlink mysql +brew install mysql-connector-c +sed -i -e 's/libs="$libs -l "/libs="$libs -lmysqlclient -lssl -lcrypto"/g' /usr/local/bin/mysql_config +pip install MySQL-python +brew unlink mysql-connector-c +brew link --overwrite mysql +``` + + + +## 更新PYPI源的地址 + +> vim myproject/Pipfile + +```ini +[[source]] +name = "pypi" +url = "https://pypi.doubanio.com/simple" +verify_ssl = true +``` + +* 可改为国内[豆瓣源](https://pypi.doubanio.com/simple/),加快下载速度 + +## 利用虚拟环境开发 + +```bash +pipenv shell +``` + +* 通过如上指令即可进入virtualenv虚拟开发环境,此虚拟环境默认Python Shll为3.6.6 + +# 错误汇总 + +```bash +# 错误详情 +--- +Uninstalling setuptools-18.5: +Could not install packages due to an EnvironmentError: [('/System/Library/Frameworks/Python.framework/Versions/2.7/Extras/lib/python/_markerlib/markers.pyc', '/private/tmp/pip-uninstall-xG0njw/System/Library/Frameworks/Python.framework/Versions/2.7/Extras/lib/python/_markerlib/markers.pyc', "[Errno 1] Operation not permitted: '/private/tmp/pip-uninstall-xG0njw/System/Library/Frameworks/Python.framework/Versions/2.7/Extras/lib/python/_markerlib/markers.pyc'"), ('/System/Library/Frameworks/Python.framework/Versions/2.7/Extras/lib/python/_markerlib/__init__.py', '/private/tmp/pip-uninstall-xG0njw/System/Library/Frameworks/Python.framework/Versions/2.7/Extras/lib/python/_markerlib/__init__.py', "[Errno 1] Operation not permitted: '/private/tmp/pip-uninstall-xG0njw/System/Library/Frameworks/Python.framework/Versions/2.7/Extras/lib/python/_markerlib/__init__.py'"), ('/System/Library/Frameworks/Python.framework/Versions/2.7/Extras/lib/python/_markerlib/markers.py', '/private/tmp/pip-uninstall-xG0njw/System/Library/Frameworks/Python.framework/Versions/2.7/Extras/lib/python/_markerlib/markers.py', "[Errno 1] Operation not permitted: '/private/tmp/pip-uninstall-xG0njw/System/Library/Frameworks/Python.framework/Versions/2.7/Extras/lib/python/_markerlib/markers.py'"), ('/System/Library/Frameworks/Python.framework/Versions/2.7/Extras/lib/python/_markerlib/__init__.pyc', '/private/tmp/pip-uninstall-xG0njw/System/Library/Frameworks/Python.framework/Versions/2.7/Extras/lib/python/_markerlib/__init__.pyc', "[Errno 1] Operation not permitted: '/private/tmp/pip-uninstall-xG0njw/System/Library/Frameworks/Python.framework/Versions/2.7/Extras/lib/python/_markerlib/__init__.pyc'"), ('/System/Library/Frameworks/Python.framework/Versions/2.7/Extras/lib/python/_markerlib', '/private/tmp/pip-uninstall-xG0njw/System/Library/Frameworks/Python.framework/Versions/2.7/Extras/lib/python/_markerlib', "[Errno 1] Operation not permitted: '/private/tmp/pip-uninstall-xG0njw/System/Library/Frameworks/Python.framework/Versions/2.7/Extras/lib/python/_markerlib'")] +--- +# 解决办法 +--- +pip install pipenv --upgrade --ignore-installed +--- +``` + diff --git "a/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\231\232\346\213\237\347\216\257\345\242\203/pyenv.md" "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\231\232\346\213\237\347\216\257\345\242\203/pyenv.md" new file mode 100644 index 0000000..17fb03d --- /dev/null +++ "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\231\232\346\213\237\347\216\257\345\242\203/pyenv.md" @@ -0,0 +1,104 @@ + + +--- + +* [模块简介](#模块简介) +* [实现原理](#实现原理) +* [调用顺序](#调用顺序) +* [安装部署](#安装部署) +* [常规套路](#常规套路) + * [查看可安装版本](#查看可安装版本) + * [安装指定的版本](#安装指定的版本) + * [查看已安装版本](#查看已安装版本) + * [卸载指定的版本](#卸载指定的版本) + * [环境变量切换版本](#环境变量切换版本) + * [删除版本环境变量](#删除版本环境变量) + * [本地切换解释器](#本地切换解释器) + * [全局切换解释器](#全局切换解释器) +* [自动初始化](#自动初始化) + +---- + +# 模块简介 + +> 第三方模块,可实现一台主机上轻松安装管理切换多个Python解释器版本 + +# 实现原理 + +> 借助PATH搜索优先级,将生成的\~/.pyenv/shims插入PATH的头部,调用并执行shims目录下的垫片程序 + +# 调用顺序 + +* 查看PYENV_VERSION环境变量是否设置,可通过pyenv shell指定或删除 +* 查看当前目录是否存在.python-version文件,可通过pyenv local设置或切换 +* 查看\~/.pyenv/下是否存在version文件,可通过pyenv global设置或切换 + +# 安装部署 + +```bash +pip install pyenv +``` + +# 常规套路 + +## 查看可安装版本 + +```bash +pyenv install --list +``` + +## 安装指定的版本 + +```bash +pyenv install 3.6.6 +pyenv rehash +``` + +## 查看已安装版本 + +```bash +pyenv versions +``` + +## 卸载指定的版本 + +```bash +pyenv uninstall 3.7.0 +pyenv rehash +``` + +## 环境变量切换版本 + +```python +pyenv shell 3.6.6 +echo ${PYENV_VERSION} +``` + +## 删除版本环境变量 + +```bash +pyenv shell --unset +``` + +## 本地切换解释器 + +```bash +pyenv local 3.6.6 +``` + +* 设置的值将写入./python-version中,下次启动时shims垫片程序会尝试读取此文件并启动指定解释器 + +## 全局切换解释器 + +```bash +pyenv global system +``` + +* 设置的值将写入~/.pyenv/version中,下次启动时shims垫片程序会尝试读取此文件并启动指定解释器 + +# 自动初始化 + +```bash +echo -e '# pyenv\nif command -v pyenv 1>/dev/null 2>&1; then\n eval "$(pyenv init -)"\nfi' >> ~/.bash_profile +``` + diff --git "a/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/asyncio.md" "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/asyncio.md" new file mode 100644 index 0000000..1777f86 --- /dev/null +++ "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/asyncio.md" @@ -0,0 +1,19 @@ +# 模块由来 + +> GIL锁的存在,多线程环境下IO密集型任务的瓶颈在于上下文切换和线程锁竞争,所以引入协程的概念,协程的上下文切换由程序控制而非内核,相当于单线程,不存在锁的概念,所以速度和性能上都比线程更有优势. + +# 模块简介 + +> 内置模块,依赖事件循环调度器,将事件生成器显式注册到调度器上,调度器对生成器进行循环调用, + +## 模拟实现 + +```python + +``` + + + +# 环境依赖 + +> \>=Python3 \ No newline at end of file diff --git "a/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/multiprocessing.assets/image-20181213161310651-4688790.png" "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/multiprocessing.assets/image-20181213161310651-4688790.png" new file mode 100644 index 0000000..982e400 Binary files /dev/null and "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/multiprocessing.assets/image-20181213161310651-4688790.png" differ diff --git "a/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/multiprocessing.assets/image-20181213161310651.png" "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/multiprocessing.assets/image-20181213161310651.png" new file mode 100644 index 0000000..982e400 Binary files /dev/null and "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/multiprocessing.assets/image-20181213161310651.png" differ diff --git "a/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/multiprocessing.assets/image-20181215082826704-4833707.png" "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/multiprocessing.assets/image-20181215082826704-4833707.png" new file mode 100644 index 0000000..c350b7c Binary files /dev/null and "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/multiprocessing.assets/image-20181215082826704-4833707.png" differ diff --git "a/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/multiprocessing.assets/image-20181215082826704.png" "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/multiprocessing.assets/image-20181215082826704.png" new file mode 100644 index 0000000..c350b7c Binary files /dev/null and "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/multiprocessing.assets/image-20181215082826704.png" differ diff --git "a/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/multiprocessing.md" "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/multiprocessing.md" new file mode 100644 index 0000000..cf5ffee --- /dev/null +++ "b/Python/\345\270\270\347\224\250\346\250\241\345\235\227/\350\277\233\347\272\277\345\215\217\347\250\213/multiprocessing.md" @@ -0,0 +1,1289 @@ +---- + +* [模块由来](#模块由来) + +* [简单介绍](#简单介绍) + +* [创建进程](#创建进程) + + * [创建方式](#创建方式) + * [属性方法](#属性方法) + +* [创建进程池](#创建进程池) + + * [模拟实现](#模拟实现) + * [真实架构](#真实架构) + * [创建方式](#创建方式) + * [属性方法](#属性方法) + +* [创建线程](#创建线程) + +* [进程间资源共享](#进程间资源共享) + + * [Pipe](#Pipe) + * [Queue](#Queue) + * [Value](#Value) + * [Array](#Array) + * [Manager](#Manager) + +* [进程间竞态同步](#进程间竞态同步) + + * [非进程池](#非进程池) + + * [Lock](#Lock) + * [RLock](#RLock) + * [Semaphore](#Semaphore) + * [Event](Event) + * [Condition](#Condition) + + * [进程池](#进程池) + + * [巧用initializer和initargs](#巧用initializer和initargs) + + * [分布式应用支持](#分布式应用支持) + + * [常规套路](#常规套路) + + * [踩坑日记](#踩坑日记) + +--- + +# 模块由来 + +> Python默认使用全局解释锁GIL保护解释器内置状态数据,例如引用计数,导致同一时间只能有一个线程操作解释器,所以所谓的多线程只能利用单核,为了充分利用多核,通常会配合多进程 + +# 简单介绍 + +> 内置模块,简化多进程,多线程代码编写,支持进线程间通信,共享,同步等特性 + +# 创建进程 + +## 创建方式 + +* 直接通过multiprocessing.Process创建 + +| 方法 | 说明 | +| ------------------------------------------------------------ | ------------------------------------------------------------ | +| Process(group=None, target=None, name=None, args=(), kwargs={}) | group未实现可忽略,target为目标函数或方法,name为进程名,args为target的位置参数,kwargs为target为命名参数 | + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import multiprocessing + + +def execute_cmd(cmd, **kwargs): + # 获取当前运行的进程对象 + p = multiprocessing.current_process() + print 'Process({0}#{1}) execute cmd: {2}'.format(p.name, p.pid, cmd) + + +if __name__ == '__main__': + # 创建一个 + # 进程名为cmd-executor + # 目标函数为execute_cmd + # 目标函数位置参数为('pwd',) + # 目标函数命名参数为{} + # 子进程对象 + p = multiprocessing.Process(name='cmd-executor', target=execute_cmd, args=('pwd',), kwargs={}) + # 启动子进程 + p.start() + # 等待子进程结束后再继续执行下面的代码 + p.join() +``` + +* 继承multiprocessing.Process子类创建 + +````python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import multiprocessing + + +class ExecuteCmdProcess(multiprocessing.Process): + def __init__(self, **kwargs): + super(ExecuteCmdProcess, self).__init__(**kwargs) + # 获取cmd + self._cmd = self._args[0] + + # 其实通过multiprocessing.Process时候调用的即此方法 + def run(self): + print 'Process({0}#{1}) execute cmd: {2}'.format(self.name, self.pid, self._cmd) + + +if __name__ == '__main__': + # 创建一个 + # 进程名为cmd-executor + # 目标函数位置参数为('pwd',) + # 目标函数命名参数为{} + # ExecuteCmdProcess子进程对象 + p = ExecuteCmdProcess(name='cmd-executor', args=('pwd',), kwargs={}) + # 启动子进程 + p.start() + # 等待子进程结束后再继续执行下面的代码 + p.join() +```` + + + +## 属性方法 + +| 属性 | 说明 | +| -------------------- | ------------------------------------------------------------ | +| p.authkey | 分布式调度时C/S间认证密钥 | +| p.daemon | 当父进程结束时子进程是否强制结束 | +| p.exitcode | 进程运行状态返回None,执行成功返回0,被N信号中断返回-N | +| p.name | 进程名 | +| p.pid | 进程ID | +| p.is_alive() | 判断进程是否存活 | +| p.join(timeout=None) | 阻塞当前进程,直到进程执行完毕或终止或者timeout超时 | +| p.start() | 将进程加入调度队列,等待CPU调度 | +| p.terminate() | 无论任务是否完成,立即停止,需要注意的是调用之后必须再次调用.join()才能真正关闭子进程 | + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import os +import multiprocessing + + +class ExecuteCmdProcess(multiprocessing.Process): + def __init__(self, **kwargs): + super(ExecuteCmdProcess, self).__init__(**kwargs) + self._cmd = self._args[0] + + def run(self): + # 模拟创建并读取100G数据,故意使其卡住 + with open('100G.dat', 'w+b') as f: + f.seek(100*pow(1024, 3)) + f.write(os.linesep) + f.seek(0) + for _ in f: + pass + +if __name__ == '__main__': + _timeout = 5 + p = ExecuteCmdProcess(name='cmd-executor', args=('cat 100G.dat',), kwargs={}) + p.daemon = True + p.start() + p.join(timeout=_timeout) + p.terminate() + # 注意调用.terminate()方法后一定要再次调用.join()才能真正关闭子进程,否则子进程此时会呈现"僵尸"状态 + p.join() + print p.is_alive() + print 'Process({0}#{1}) exit with {2} after {3} seconds'.format(p.name, p.pid, p.exitcode, _timeout) +``` + +# 创建进程池 + +> 进程的创建和销毁是非常消耗资源的,进程池不仅可以控制进程数量而且还可以控制进程创建和销毁,内部使用一个独立线程默认每隔0.1秒检测进程池内是否有需要清理的worker子进程,如果有则先清理然后在创建同等数量的新的worker子进程加入进程池,其实每个worker的主要任务就是每次从任务队列中取出任务执行再将结果写入结果队列,所以可以很容易控制worker子进程最大执行的任务数,超出则结束此worker子进程,独立线程检测到后立即清理并创建新的子进程 + +## 模拟实现 + +```python + +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import time +import threading +import multiprocessing + + +# 工作进程目标函数 +# 1. 从任务队列读取任务 +# 2. 执行任务并将结果写入结果队列 +def worker(qtask, qresult): + while True: + try: + task = qtask.get() + if task is None: + print 'Recv stop signal from workers_check thread, gracefull exit' + break + func, args, kwargs, callback = task + except Exception as e: + print 'Get task from task_queue with exception, exp={0}'.format(e) + break + + try: + result = (True, func(*args, **kwargs)) + except Exception as e: + result = (False, e) + callable(callback) and callback(result) + + try: + qresult.put(result) + except Exception as e: + print 'Put result to result queue with exception, exp={0}'.format(e) + + +# 进程池实现 +class Pool(object): + Process = multiprocessing.Process + + def __init__(self, max_processes=None): + # 进程池最大进程数 + self._max_processes = max_processes or multiprocessing.cpu_count() + + # 进程池 + self._pool = [] + # 任务队列 + self._task_queue = multiprocessing.Queue() + # 结果队列 + self._rest_queue = multiprocessing.Queue() + + self._workers_check = threading.Thread( + target=Pool._check_workers, + args=(self,) + ) + + # 模拟优雅退出,零时方案 + def start_workers_check(self): + self._workers_check.daemon = True + self._workers_check.start() + + @staticmethod + def _check_workers(p): + # 定期检查是否需要创建新的工作进程并创建 + while not p._task_queue.empty(): + print 'Task queue is not empty, keep checking' + p._keep_workers() + time.sleep(0.1) + # 否则发送None使工作进程停止 + print 'Task queue is empty, notify all worker processes' + map(lambda _: p._task_queue.put(None), p._pool) + + def _need_clean(self): + clean = not self._pool + for i in reversed(range(len(self._pool))): + w = self._pool[i] + # exitcode为None表示运行中 + if w.exitcode is not None: + w.join() + clean = True + print 'Clean exit worker {0}'.format(w) + # 清理已停止的进程 + del self._pool[i] + return clean + + def _fill_workers(self): + need_worker_num = self._max_processes - len(self._pool) + for _ in xrange(need_worker_num): + # 创建新的工作进程 + p = self.Process( + target=worker, + args=(self._task_queue, self._rest_queue) + ) + p.name = p.name.replace('Process', '{0}{1}'.format(self.__class__.__name__, p.name)) + p.daemon = True + # 加入进程池 + self._pool.append(p) + p.start() + print '{0} worker processes add to pool'.format(need_worker_num) + + def _keep_workers(self): + self._need_clean() and self._fill_workers() + + def map(self, func, iterable, callback=None): + map(lambda _: self._task_queue.put((func, (_,), {}, callback)), iterable) + + def join(self): + self.start_workers_check() + print 'Worker process check thread start' + # 等待工作进程检查线程结束 + self._workers_check.join() + # 等待工作进程关闭 + for p in self._pool: + p.terminate() + p.join() + + +def execute_cmd(cmd): + p = multiprocessing.current_process() + print '{0}#{1} execute cmd succ, cmd={2}'.format(p.name, p.pid, cmd) + + return cmd + + +def execute_cmd_callback(result): + print 'Got result {0}'.format(result) + +if __name__ == '__main__': + p = Pool() + # 模拟执行1万个命令执行 + p.map(execute_cmd, xrange(10000), execute_cmd_callback) + p.join() +``` + +* 为了加深大家对进程池概念的理解,如上简单模拟进程池实现,但真正的进程池远不止这么简单,如下为完整内部架构图 + +## 真实架构 + +![image-20181213161310651](multiprocessing.assets/image-20181213161310651-4688790.png) + +## 创建方式 + +| 方法 | 说明 | +| ------------------------------------------------------------ | ------------------------------------------------------------ | +| Pool(processes=None, initializer=None, initargs=(), maxtasksperchild=None) | processes为进程数,如果为None则默认为multiprocessing.cpu_count(),initializer和initargs忽略,maxtasksperchild为子进程生命周期内最大任务数,主要用于释放闲置资源,如果为None则表示进程一直存活 | + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import multiprocessing + + +def execute_cmd(cmd): + p = multiprocessing.current_process() + print '{0}#{1} execute cmd succ, cmd={2}'.format(p.name, p.pid, cmd) + + return cmd + + +if __name__ == '__main__': + p = multiprocessing.Pool() + print p.map(execute_cmd, xrange(10000)) +``` + +## 属性方法 + +| 方法 | 说明 | +| ---------------------------------------------------------- | ------------------------------------------------------------ | +| p.apply(func, args=(), kwds={}) | 创建一个任务给进程池,阻塞模式,直接返回结果,func表示目标函数或方法,args表示func的位置参数,kwds表示func的命名参数 | +| p.apply_async(func, args=(), kwds={}, callback=None) | 创建一个任务给进程池,非阻塞模式,返回ApplyResult对象,通过其get方法获取返回值,func表示目标函数或方法,args表示func的位置参数,kwds表示func的命名参数,callback表示回调函数或方法,用于异步回调模式 | +| p.map(func, iterable, chunksize=None) | 分组创建多个任务,阻塞模式,多进程无序执行但顺序返回,func表示目标函数或方法,iterable表示可迭代对象,chunksize表示分组创建任务,可忽略 | +| p.map_async(func, iterable, chunksize=None, callback=None) | 分组创建多个任务,非阻塞模式,多进程无序执行但顺序返回,func表示目标函数或方法,iterable表示可迭代对象,chunksize表示分组创建任务,可忽略,callback表示回调函数或方法,用于异步回调模式 | +| p.close() | 通过状态标志位_state声明不再处理新的任务 | +| p.terminate() | 同上,优雅关闭所有涉及的线程和子进程 | +| p.join() | 等待相关进线协程都结束后再退出,调用它前需要先调用p.close() | + +* apply和apply_async的实现相对简单,首先创建一个ApplyResult结果对象,并携带结果对象序号_job写入\_taskqueue任务队列,\_handle_tasks任务处理器从任务队列取出放入\_inqueue进程处理队列,工作进程从进程处理队列中取出执行写入\_outqueue结果队列,\_handle_results结果处理器通过序号\_job找到对应的ApplyResult结果对象将结果\_set到对象 +* map和map_async在上面的基础上在创建任务的时候_get_tasks配合chunksize对任务分组,工作进程收到的任务是分组形式的任务,然后子进程内通过mapstar(其实就是巧妙的调用了map)使一个工作进程可以按组执行任务 + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import multiprocessing + + +def execute_cmd(cmd): + p = multiprocessing.current_process() + print '{0}#{1} execute cmd succ, cmd={2}'.format(p.name, p.pid, cmd) + + return cmd + + +if __name__ == '__main__': + p = multiprocessing.Pool(processes=multiprocessing.cpu_count()) + for i in xrange(1000): + # 同步写入 + p.apply_async(execute_cmd, args=(i,)) + + # 必须先close设置标志位不再接收新的任务,或者保证主进程不退出也ok + p.close() + p.join() +``` + +# 创建线程 + +> multiprocessing.dummy + +# 进程间资源共享 + +> 模块内置Pipe,Queue,Value, Array,Manager等可轻松实现进程间资源共享,但需要注意的是由于底层IPC或RPC通信序列化类只支持Pickle和Xmlrpclib,而默认的Pickle不支持对复杂对象序列化,即使强制支持了序列化后期还原对象时也可能由于魔术方法而引发其它问题 + +### Pipe + +> 可读可写,但只适用于Process类,不能用于Pool类,而且只能用于两个进程间全双工或半双工通信 + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import itertools +import multiprocessing + + +job_counter = itertools.count() + + +# 生产者 +def producer(p): + while True: + p.send(job_counter.next()) + + +# 消费者 +def consumer(p): + while True: + job = p.recv() + print 'Recv job#{0}'.format(job) + + +if __name__ == '__main__': + # 全双工通道 + r, w = multiprocessing.Pipe(True) + # 消费者进程 + c = multiprocessing.Process(target=consumer, name=consumer.__name__, args=(r,)) + # 生产者进程 + p = multiprocessing.Process(target=producer, name=producer.__name__, args=(w,)) + c.daemon = True + p.daemon = True + c.start() + p.start() + c.join() + p.join() +``` + +### Queue + +> 可读可写,但只适用于Process类,不能用于Pool类,可用于多个进程间共享 + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import itertools +import multiprocessing + + +job_counter = itertools.count() + +# 生产者 +def producer(queue): + while True: + queue.put(job_counter.next()) + +# 消费者 +def consumer(queue): + while True: + job = queue.get() + print 'Recv job#{0}'.format(job) + +if __name__ == '__main__': + # 共享队列 + q = multiprocessing.Queue() + # 消费者进程 + c = multiprocessing.Process(target=consumer, name=consumer.__name__, args=(q,)) + # 生产者进程 + p = multiprocessing.Process(target=producer, name=producer.__name__, args=(q,)) + c.daemon = True + p.daemon = True + c.start() + p.start() + c.join() + p.join() +``` + +* Queue对象还支持put, get, qsize, empty, full等常用方法,具体用法请深入源码 + +### Value + +### Array + +> 可读可写,但只适用于Process类,不能用于Pool类,可用于多个进程间共享 + +* 思考 + * 如下代码为何无法实现简单的生产者消费者模型 ? 问题出在哪里 ? + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import time +import itertools +import multiprocessing + + +job_counter = itertools.count() + + +# 生产者 +def producer(share_list): + while True: + share_list.append(job_counter.next()) + + +# 消费者 +def consumer(share_list): + while True: + if not share_list: + print 'Empty list, ignore' + time.sleep(0) + continue + job = share_list.pop() + print 'Recv job#{0}'.format(job) + + +if __name__ == '__main__': + # "共享"内存 + share_list = [] + # 消费者进程 + c = multiprocessing.Process(target=consumer, name=consumer.__name__, args=(share_list,)) + # 生产者进程 + p = multiprocessing.Process(target=producer, name=producer.__name__, args=(share_list,)) + c.daemon = True + p.daemon = True + c.start() + p.start() + c.join() + p.join() +``` + +* 由于在Fork子进程时其实对主进程的上下文包括环境变量都会复制一份儿给子进程,而此时复制的新的变量并不是指向原来share_list的内存地址,所以无法同步修改 + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import itertools +import multiprocessing + + +job_counter = itertools.count() + + +# 生产者 +def producer(share_temp, share_list): + share_list_len = len(share_list) + while True: + # 如果单值共享内存值为0则创建10个计数写入 + if share_temp.value == 0: + for i in xrange(share_list_len): + share_list[i] = job_counter.next() + # 模拟信号告诉消费者消费 + share_temp.value = 1 + + +# 消费者 +def consumer(share_temp, share_list): + while True: + # 如果单值共享内存值为1则尝试读取所有的值 + if share_temp.value == 1: + for job in share_list: + print 'Recv job#{0}'.format(job) + # 模拟信号告知生产者消费完毕 + share_temp.value = 0 + +if __name__ == '__main__': + # 共享内存 + # multiprocessing.sharedctypes.typecode_to_type + mem_len = 10 + share_temp = multiprocessing.Value('i', 0) + share_list = multiprocessing.Array('i', [0]*mem_len) + # 消费者进程 + c = multiprocessing.Process(target=consumer, name=consumer.__name__, args=(share_temp, share_list,)) + # 生产者进程 + p = multiprocessing.Process(target=producer, name=producer.__name__, args=(share_temp, share_list,)) + c.daemon = True + p.daemon = True + c.start() + p.start() + c.join() + p.join() +``` + +* Value和Array都支持进程间共享,区别是前者只能存储单值,后者可以存储多个值,至于类型请参考源码multiprocessing.sharedctypes.typecode_to_type + +### Manager + +> 可读可写,不仅适用于Process类,而且适用于Pool类,可用于多个进程间共享,但需要注意的是Win下Manager对象必须在\_\_main\_\_下声明 + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import itertools +import multiprocessing + + +job_counter = itertools.count() + + +# 生产者 +def producer(share_list): + while True: + share_list.append(job_counter.next()) + + +# 消费者 +def consumer(share_list): + while True: + try: + job = share_list.pop() + except IndexError: + continue + print 'Recv job#{0}'.format(job) + + +if __name__ == '__main__': + # 共享内存 + share_list = multiprocessing.Manager().list() + # 消费者进程 + c = multiprocessing.Process(target=consumer, name=consumer.__name__, args=(share_list,)) + # 生产者进程 + p = multiprocessing.Process(target=producer, name=producer.__name__, args=(share_list,)) + c.daemon = True + p.daemon = True + c.start() + p.start() + c.join() + p.join() +``` + +* Manager其实调用的multiprocessing.managers.SyncManager,它支持很多其它数据类型,具体可深入源码 + +# 进程间竞态同步 + +> 多个进程同时读写资源可能由于竞态而导致资源数据混乱或死锁 + +## 非进程池 + +### Lock + +> 互斥锁,支持上下文管理器,作为公共锁,一旦一个进程获得锁其它进程再尝试获取将阻塞直至此锁被释放 + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import time +import multiprocessing + + +class Account(object): + def __init__(self, name, money): + self.name = name + self.lock = multiprocessing.Lock() + self.money = money + + def put(self, money): + self.money += money + + def pay(self, money): + m = self.money - money + self.money = m if m > 0 else 0 + + +def transfer(f, t, m): + # 首先锁住转钱者,可能是user1,也可能是user2 + print 'Account {0} require lock'.format(t.name) + with f.lock: + print 'Account {0} lock succ'.format(f.name) + # 转钱 + f.pay(m) + print 'Account {0} transfer {1} to {2}'.format(f.name, m, t.name) + # 交出控制权 + time.sleep(1) + # 尝试锁住收钱者,但由于CPU切换可能 + print 'Account {0} require lock'.format(t.name) + with t.lock: + print 'Account {0} lock succ'.format(f.name) + # 存钱 + t.put(m) + print 'Account {0} got {1} from {2}'.format(f.name, m, t.name) + print 'Account {0} release lock'.format(f.name) + print 'Account {0} release lock'.format(t.name) + +if __name__ == '__main__': + # 分别创建user1和user2 + user1 = Account('user1', 1000) + user2 = Account('user2', 2000) + + # 多进程模式下user1向user2转100元,与此同时user2向user1转200,所以必须加锁 + map(lambda p: p.start(), [ + multiprocessing.Process(target=transfer, args=(user1, user2, 100)), + multiprocessing.Process(target=transfer, args=(user2, user1, 200)) + ]) +``` + +* 如上简单模拟转钱时锁定账户导致的死锁,负责user1转钱给user2的进程和user2转钱给user1的进程由于是并行运行,所以都会首先抢占一个锁,到后面锁定对方账户的时候由于对方账户的锁都未释放,所以会出现相互等待的状态,也就是死锁状态 + +### RLock + +> 递归锁,支持上下文管理器,作为公共锁,一个进程可以多次获得锁,内部会为其维护一个独立counter,记录获得锁的次数,但需要注意的是只有此进程释放了counter锁之后其它进程才可以操作锁住的资源 + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import time +import multiprocessing + + +class Account(object): + lock = multiprocessing.RLock() + + def __init__(self, name, money): + self.name = name + self.money = money + + def put(self, money): + self.money += money + + def pay(self, money): + m = self.money - money + self.money = m if m > 0 else 0 + + +def transfer(f, t, m): + # 首先锁住转钱者,可能是user1,也可能是user2 + print 'Account {0} require lock'.format(t.name) + f.lock.acquire() + print 'Account {0} lock succ'.format(f.name) + # 转钱 + f.pay(m) + print 'Account {0} transfer {1} to {2}'.format(f.name, m, t.name) + # 交出控制权 + time.sleep(1) + # 尝试锁住收钱者,但由于CPU切换可能 + print 'Account {0} require lock'.format(t.name) + t.lock.acquire() + print 'Account {0} lock succ'.format(f.name) + # 存钱 + t.put(m) + print 'Account {0} got {1} from {2}'.format(f.name, m, t.name) + # 释放其它的锁 + t.lock.release() + # 释放自己的锁 + f.lock.release() + print 'Account {0} release lock'.format(f.name) + print 'Account {0} release lock'.format(t.name) + +if __name__ == '__main__': + # 分别创建user1和user2 + user1 = Account('user1', 1000) + user2 = Account('user2', 2000) + + # 多进程模式下user1向user2转100元,与此同时user2向user1转200,所以必须加锁 + map(lambda p: p.start(), [ + multiprocessing.Process(target=transfer, args=(user1, user2, 100)), + multiprocessing.Process(target=transfer, args=(user2, user1, 200)) + ]) +``` + +### Semaphore + +> 信号量,支持上下文管理器,作为公共锁池,多个进程竞态获取锁池内的锁,如果锁的数量为0则其它进程等待锁被释放,其实内部维护一个公共的counter,限制并发 + +### Event + +> 事件,支持上下文管理器,作为公共标志位,当一个或多个进程依赖另一个进程状态时可通过set()设置全局标志位,clear()清空全局标志位,is_set()查看是否全局标志位被设置,wait()阻塞等待全局标志位被设置 + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import time +import MySQLdb +import settings +import itertools +import threading +from Queue import Queue +from contextlib import contextmanager +from MySQLdb.cursors import DictCursor + + +conn_counter = itertools.count() + + +# 连接代理对象 +class ConnProxy(object): + def __init__(self, inst, conn): + self._inst = inst + # 真正的基于MySQLdb创建的连接 + self._conn = conn + # 连接ID,主要为了记录 + self._conn_id = conn_counter.next() + + @property + def c_id(self): + return self._conn_id + + @property + def real(self): + return self._conn + + def release(self): + self._inst.release(self) + + +# 数据库连接池对象 +class DBPool(object): + def __init__(self, *args, **kwargs): + self._args = args + self._kwargs = kwargs + self._conn_pool_size = kwargs.get('conn_pool_size', 10) + self._conn_pool_queue = Queue(maxsize=self._conn_pool_size) + # 连接池并发请求数限制,演示而已,因为Queue会自动阻塞其它线程 + self._semaphore = threading.Semaphore(kwargs.get('concurrent_conns', 5)) + + self._init_pool_queue() + + # 连接池保持线程 + conn_pool_keeper = threading.Thread( + target=DBPool._conn_pool_keeper, + args=(self,) + ) + conn_pool_keeper.setDaemon(True) + conn_pool_keeper.start() + + @property + def conn_pool_size(self): + return self._conn_pool_size + + @property + def conn_pool_queue_size(self): + return self._conn_pool_queue.qsize() + + @staticmethod + def _conn_pool_keeper(db): + while True: + # 周期遍历踢出已关闭的代理连接并创建新的代理连接放入连接池 + conn_proxy = db._conn_pool_queue.get() + if conn_proxy is None: + break + if conn_proxy.real.closed: + conn_proxy = db._create_new_conn() + db._conn_pool_queue.put(conn_proxy) + + time.sleep(0.1) + + def _create_new_conn(self): + # 创建新的连接并封装为代理连接对象 + conn = MySQLdb.connect(*self._args, **self._kwargs) + conn_proxy = ConnProxy(self, conn) + return conn_proxy + + def _init_pool_queue(self): + # 初始化队列,创建conn_pool_size个代理连接对象放入队列 + for _ in xrange(self._conn_pool_size): + conn = self._create_new_conn() + self._conn_pool_queue.put(conn) + + # 主要配合上下文管理器自动回收代理连接对象 + @contextmanager + def acquire(self): + conn_proxy = None + # 从连接池队列获取代理连接对象 + try: + with self._semaphore: + conn_proxy = self._conn_pool_queue.get() + yield conn_proxy + finally: + self.release(conn_proxy) + + def release(self, conn_proxy): + # 将代理连接对象放回连接池 + self._conn_pool_queue.put(conn_proxy) + + +# 等待数据库连接池初始化完毕 +def db_prepare(e): + global db_pool + + db_pool = DBPool(host=settings.HOST, port=settings.PORT, + user=settings.USER, + passwd=settings.PASSWD, db=settings.DB) + while db_pool.conn_pool_size != db_pool.conn_pool_queue_size: + time.sleep(0.1) + # 通知其它线程数据库连接池对象初始化完毕 + e.set() + + +def execute_sql(e, sql, callback=None): + global db_pool + # 阻塞等待信号 + e.wait() + + # 通过上下文管理器自动获取回收关闭 + with db_pool.acquire() as db_conn: + with db_conn.real.cursor(cursorclass=DictCursor) as db_cursor: + db_cursor.execute(sql) + if callable(callback): + connid = db_conn.c_id + result = db_cursor.fetchall() + callback(connid, sql, result) + + +def execute_sql_callback(connid, sql, result): + print u'连接ID: {0} 执行语句: {1} 执行结果: {2}'.format(connid, sql, result) + + +if __name__ == '__main__': + # 全局数据库连接池对象 + db_pool = None + # 全局线程事件对象,主要用于通知其它线程数据库连接池对象准备就绪可以使用 + t_event = threading.Event() + + # 加入线程队列 + threads = [threading.Thread(target=db_prepare, args=(t_event,))] + for _ in xrange(1000): + t = threading.Thread(target=execute_sql, args=(t_event, 'SELECT VERSION() AS VERSION;', execute_sql_callback)) + threads.append(t) + t.setDaemon(True) + # 启动运行所有线程 + for t in threads: + t.start() + t.join() +``` + +* 由于多进程资源共享无法共享复杂对象,例如如上定义的数据库连接池,因为内部进程IPC/RPC通信传递的是pickle序列化后的数据,复杂对象序列化还原时会由于各种魔术方法等因素导致出错,所以如上使用多线程模拟实现数据库连接池,其实用法和多进程一样,帮助大家理解一下常用于并发限制的Semaphore信号量和Event事件的使用 + +### Condition + +> 条件锁,基于Lock/RLock的高级锁(同样的方法调用),支持上下文管理器,除此之外,.wait(timeout=None)会释放内部所有Lock/RLock,并将进程挂起,直到被.c.notify()/.c.notify_all()通知再恢复 + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import time +import random +import itertools +import multiprocessing + + +# 公共基类 +class BaseProcess(multiprocessing.Process): + def __init__(self, share_mem, condition, maxlength, *args, **kwargs): + super(BaseProcess, self).__init__(*args, **kwargs) + # 进程间共享列表 + self._share_mem = share_mem + # 条件锁 + self._condition = condition + # 列表最大长度 + self._maxlength = maxlength + + +# 生产者进程 +class Producer(BaseProcess): + def run(self): + counter = itertools.count() + while True: + with self._condition: + # 当共享列表满则释放锁,阻塞当前状态 + if len(self._share_mem) == self._maxlength: + print 'Share mem is full, wait...' + self._condition.wait() + print 'Space in share_mem and consumer notify me' + data = counter.next() + print 'Product {0}'.format(data) + self._share_mem.append(data) + # 为了不让共享列表满,一旦共享列表有数据就通知Consumer去消费,如果注释掉则可能由于Producer超出共享列表长度限制和Consumer首次判断共享列表为空而都被wait + self._condition.notify() + + +# 消费者进程 +class Consumer(BaseProcess): + def run(self): + while True: + with self._condition: + # 当共享列表为空则释放锁,阻塞当前状态 + if len(self._share_mem) == 0: + print 'Share mem is empty, wait...' + self._condition.wait() + print 'Producer put something to share_mem and notify me' + data = self._share_mem.pop() + print 'Consume {0}'.format(data) + # 为了不让共享列表空,一旦有消费就通知Producer去生产,如果注释掉可能由于Consumer消费到列表为空,而Producer依然还不知情不生产 + self._condition.notify() + + time.sleep(random.random()) + +if __name__ == '__main__': + share_mem = multiprocessing.Manager().list() + condition = multiprocessing.Condition() + + # 启动进程 + c = Consumer(share_mem, condition, 10) + p = Producer(share_mem, condition, 10) + c.start() + p.start() + c.join() + p.join() +``` + + + +## 进程池 + +> 底层IPC或RPC通信序列化类只支持Pickle和Xmlrpclib,虽然可以通过使用Manager注册的Lock,RLock,Semaphone,Event,Condition来直接在进程池使用,但其实是会启动一个独立的进程来与其它子进程进行同步通信,其实可以巧用initializer和initargs实现各种轻量级简单的锁机制 + +### 巧用initializer和initargs + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import os +import sys +import time +import requests +import functools +import multiprocessing + + +def api_optimization(baseurl, novel, output=None): + start_time = time.time() + url = '{0}{1}'.format(baseurl, novel) + requests.get(url).json() + end_time = time.time() + log = 'Request({0}) cost {1} seconds{2}'.format(novel, end_time-start_time, os.linesep) + # 尝试请求锁,此处无法直接使用with语句,如果需要可以封装一层 + lock.acquire() + if output is None: + sys.stdout.write(log) + else: + with open(output, 'a+b') as f: + f.write(log) + # 释放锁 + lock.release() + return log + + +def lock_initializer(l): + global lock + lock = l + +if __name__ == '__main__': + lock = None + l = multiprocessing.Lock() + # 在初始化时通过全局变量赋值的方式将全局锁传递进去 + p = multiprocessing.Pool(initializer=lock_initializer, initargs=(l,)) + novels = ['盗墓笔记', '完美世界', '斗罗大陆', '舞动乾坤'] * 1000 + baseurl = 'https://www.apiopen.top/novelSearchApi?name=' + logfile = 'record.log' + # 为了方便调用用functools.partial装饰一下 + api_optimization_partial = functools.partial(api_optimization, baseurl, output=logfile) + print p.map(api_optimization_partial, novels) +``` + + + +# 分布式应用支持 + +> multiprocessing.managers.BaseManager支持将进程分布在多台机器上进行分布式调度,其实原理很简单,manager启动时将读写Queue注册到网络上,worker注册上来时通过名称和authkey关联对应的读写队列,后续的数据收发就全部走网络 + +## 常规套路 + +![image-20181215082826704](multiprocessing.assets/image-20181215082826704-4833707.png) + +* 服务端 + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import Queue +import string +from multiprocessing.managers import BaseManager + + +class Server(object): + task_func_name = 'get_task_queue' + result_func_name = 'get_result_queue' + + def __init__(self, manager_class, *args, **kwargs): + self._args = args + self._state = None + self._kwargs = kwargs + self._manager = None + self._manager_class = manager_class + + @property + def manager_class(self): + if self._manager_class is not BaseManager and \ + issubclass(self._manager_class, BaseManager): + return self._manager_class + # 向QueueManager注册一个任务队列和结果队列,它会自动暴露在网络中 + qtask = Queue.Queue() + qresult = Queue.Queue() + self._manager_class = type('QueueManager', (BaseManager,), {}) + # 注册后必须通过实例的.get_task_queue和.get_result_queue获取包装后的队列,不能直接操作源Queue + self.manager_class.register(self.task_func_name, callable=lambda: qtask) + self.manager_class.register(self.result_func_name, callable=lambda: qresult) + + return self._manager_class + + @property + def manager(self): + if isinstance(self._manager, self.manager_class): + return self._manager + # 创建一个管理器对象 + self._manager = self.manager_class(*self._args, **self._kwargs) + + return self._manager + + # 向任务队列写入任务 + def put_data(self, *args, **kwargs): + # 通过.get_task_queue获取任务队列代理对象 + q = getattr(self.manager, self.task_func_name)() + q.put(*args, **kwargs) + + # 从结果队列获取结果 + def get_data(self, *args, **kwargs): + # 通过.get_result_queue获取结果队列代理对象 + q = getattr(self.manager, self.result_func_name)() + return q.get(*args, **kwargs) + + def __enter__(self): + return self + + def __exit__(self, exc_type, exc_val, exc_tb): + # 关闭管理器 + self.manager.shutdown() + + +if __name__ == '__main__': + # authkey为服务端认证客户端的密钥 + with Server(BaseManager, address=('', 65533), authkey='secret') as server: + # 启动管理器 + server.manager.start() + for s in string.letters: + print 'Server send data: {0}'.format(s) + server.put_data(s) + print 'Wait client data...' + while True: + r = server.get_data() + print 'Server Recv data: {0}'.format(r) +``` + +* 客户端 + +```python +#! -*- coding: utf-8 -*- + + +# author: forcemain@163.com + + +import Queue +import string + +from multiprocessing.managers import BaseManager + + +class Client(object): + task_func_name = 'get_task_queue' + result_func_name = 'get_result_queue' + + def __init__(self, manager_class, *args, **kwargs): + self._args = args + self._kwargs = kwargs + self._manager = None + self._manager_class = manager_class + + @property + def manager_class(self): + if self._manager_class is not BaseManager and \ + issubclass(self._manager_class, BaseManager): + return self._manager_class + # 向QueueManager查找一个任务队列和结果队列 + self._manager_class = type('QueueManager', (BaseManager,), {}) + # 只是查找和映射所以无需注册时写callable + self._manager_class.register(self.task_func_name) + self._manager_class.register(self.result_func_name) + + return self._manager_class + + @property + def manager(self): + if isinstance(self._manager, BaseManager): + return self._manager + self._manager = self.manager_class(*self._args, **self._kwargs) + # 创建一个管理器对象 + return self._manager + + # 从任务队列获取任务 + def get_data(self, *args, **kwargs): + # 通过.get_task_queue获取任务队列代理对象 + q = getattr(self.manager, self.task_func_name)() + return q.get(*args, **kwargs) + + # 向结果队列写入结果 + def put_data(self, *args, **kwargs): + # 通过.get_result_queue获取结果队列代理对象 + q = getattr(self.manager, self.result_func_name)() + q.put(*args, **kwargs) + + def __enter__(self): + return self + + def __exit__(self, exc_type, exc_val, exc_tb): + pass + +if __name__ == '__main__': + # authkey为服务端认证客户端的密钥 + with Client(BaseManager, address=('127.0.0.1', 65533), authkey='secret') as client: + # 连接目标服务器 + client.manager.connect() + while True: + data = client.get_data() + print 'Client recv data: {0}'.format(data) + client.put_data(data.upper()) + print 'Client send data: {0}'.format(data.upper()) +``` + +# 踩坑日记 + +* Win平台下使用前需要调用multiprocessing.freeze_support(),否则可能会抛出RuntimeError异常 + + + diff --git "a/Python/\346\272\220\347\240\201\345\205\261\350\257\273/\345\215\263\346\227\266\346\240\207\350\256\260\345\274\225\346\223\216/README.md" "b/Python/\346\272\220\347\240\201\345\205\261\350\257\273/\345\215\263\346\227\266\346\240\207\350\256\260\345\274\225\346\223\216/README.md" new file mode 100644 index 0000000..d45cc40 --- /dev/null +++ "b/Python/\346\272\220\347\240\201\345\205\261\350\257\273/\345\215\263\346\227\266\346\240\207\350\256\260\345\274\225\346\223\216/README.md" @@ -0,0 +1,8 @@ +# 需求整理 + +> 如何将纯文本自动转换为 + +# 项目简介 + +> 简单易用的标记系统 + diff --git a/README.md b/README.md index 99812fc..5aa32bd 100644 --- a/README.md +++ b/README.md @@ -1,2 +1,6 @@ -# 读书笔记 +# 高级运维开发知识体系V1.0 + +> 作者: 李满满 +> +> 邮箱: forcemain@163.com