首页 文章 精选 留言 我的

精选列表

搜索[个人博客],共10000篇文章
优秀的个人博客,低调大师

【博客大赛】JVM深入探索之初探ZGC的领域

背景概要 虽然每一门编程语言都会有属于自己的生命周期,有巅峰期同样也有衰退期,Java语言为了保证的巅峰状态一直在线,不断的升华和创新自己特性和能力,经历了从1995年到2020年这么多年的成长,Java技术体系亦愈渐强大,无论是Java虚拟机不断的升级优化还是Java语法糖的不断完善,每一次发布都是一个非同一般的里程碑节点,要说Java体系中最重要的“大脑”,那非JVM虚拟机莫属了,而JVM中最重要的核心部分就是GC管理子系统。目前,最新颖且最优秀的GC回收器就是ZGC,对于ZGC来讲它的算法结构和设计思想非常值得大家学习和借鉴,无论你对JVM了解程度如何,相信本篇文章都会帮助对ZGC的特性和原理有一个非常深刻且扎实的认识。 基础概念 ZGC的官方定义是:"A Scalable Low Latency Garbage Collector"。翻译成中文含义就是一个可扩展的低延迟垃圾收集器。 是Oracle 内部研发的一款新的垃圾回收器,最大的宣传点是低延迟(Low Latency)。根据近期的消息,该项目将会托管给 OpenJDK,由后者开发完善。体现出:可扩展、低延迟 实现目标 ZGC特性 备注 能够支持TB级别的堆内存回收能力 10ms的最大回收停顿时间 未来回收的系统运作机制奠定了基础 与G1相比较,吞吐率下降的不多于15% 暂定时间不会随着堆活跃的对象区域的大小提升而变长 性能对比 在SPECjbb这个基准测试中,被测产品要运行JVM中包含了并行机制GC回收器(Parallel Scavenge/Parallel Old),ZGC回收器、G1回收器三者的比较。 比较的环境的CPU模型名称:Oracle Linux 7.4 - Intel Xeon ES-2690 -2.9GHz 例如:Intel(R) Xeon(R) CPU E5-2682 v4 @ 2.50GHz或者QEMU Virtual CPU version (cpu64-rhel6)应该能大概看得出是虚拟CPU仍是真实的CPU。 CPU模型名称 堆内存 操作系统(OS) 物理CPU信息 主频 物理内核 逻辑内核 Composite 128G Oracle Linux 7.4 Intel Xeon ES-2690 2.9GHz 16 32 SPECjbb®2015 – Score 比较的维度: max-JOPS(Throughput):吞吐率 ZGC Parallel G1 98% 100% 95% critical-JOPS(Throughput): 延迟角度的吞吐率 ZGC Parallel G1 76% 51% 49% 结论 总体分布数据为越高越好。 max-JOPS吞吐率:Parallel=100%,ZGC=98%,G1=95% critical-JOPS延迟角度吞吐率:ZGC=76%,Parallel=51%,G1=49% 总体来讲ZGC的综合性能最好,仅次于Parallel,最后是G1 SPECjbb®2015 – Pause Times 比较的维度: ZGC Linear Scale:线性指标 Average 95% 99% 99.9% Max 1.091ms (+/-0.215ms) 1.380ms 1.512ms 1.663ms 1.681ms Logarithmic scale :逻辑指标 Average 95% 99% 99.9% Max 0.18MS 0.5MS 0.75MS 5ms 6.1Ms 而对于其他GC回收器来讲,可以看出延迟时间大大增加,此处就不进行数据列举,参考上图可以分析全部远远>10ms之上。 结论 越低越好 各种SPEC®基准测试和内部工作负载进行了特别的性能测量。一般情况下,ZGC能够维护个位数的毫秒暂停时间。 功能特性 读屏障(Load barriers ) 染色指针(Colored pointers) 单代(不再区分新生代和老年代)(Single generation) 局部压缩(Partial compaction) 基于区域(Region-based) 即时内存重用(Immediate memory reuse) NUMA机制(NUMA-aware) 并发机制(Concurrent) 总结成为一句话就是: ZGC是一个并发的、单代(不再区分新生代和老年代)的、基于region的、支持numa的压缩收集器。Stop-the-world阶段仅限于根扫描,所以GC暂停时间不会随着堆或存活对象的多少而增加。 回收阶段 ZGC的一个核心设计原则是结合使用内存屏障和染色对象指针。 ZGC的垃圾回收算法和传统的Stop-The-World式的垃圾回收算法不太一样,后者的标记阶段和内存压缩阶段会使得应用线程挂起。ZGC 和C4(Continuously Concurrent Compacting Collector))算法比较类似。 垃圾回收过程主要分为以下三个阶段: 上图整体是GC的回收过程的周期阶段展示,主要包括了并行化处理阶段: 标记(Marking); 重定位(Relocation)/压缩(Compaction); 重新分配集的选择(Relocation set selection); 引用处理(Reference processing); 弱引用的清理(WeakRefs Cleaning); 字符串常量池(String Table)和符号表(Symbol Table)的清理; 类卸载(Class unloading)。 ZGC能够在运行Java应用程序线程时执行并发操作,比如对象重定位。从Java线程的角度来看,在Java对象中加载引用字段的行为受到内存屏障的限制。 一次完整的 ZGC 回收周期分为以下几个阶段(Phase): Pause Mark Start:标记根对象; Concurrent Mark:并发标记阶段; Concurrent Relocate:并发重定位; 活动对象被移动到了一个新的Heap RegionB-region 中,之前旧对象所在的Heap RegionA-region 即可复用;如果 B-region 中对象之间的引用关系将会在这一阶段被更新; 在重定位过程中,新旧对象的映射关系(同一对象在不同Region中的映射关系)被记录在了Forwarding Tables中。 Pause Mark Start:这个阶段实际上已经进入了新的ZGC Cycle,同样也是标记根对象; Concurrent Relocate…… 从上面的垃圾回收过程可以看到,正是因为ZGC回收过程中各个 Phase 的并发性,才使得GC Pause不受垃圾回收周期内堆上活动数据数量和需要跟踪与更新的引用数量的影响,将暂停时间保持在较低的水平。 Pause Mark Start ZGC 的垃圾回收各阶段也不都是并发执行的,在Pause Mark Start阶段进行根对象扫描(Root Scanning)时会出现短暂STW的暂停。 ZGC 在Pause Mark Start阶段进行根对象扫描(Root Scanning)建立起到达性分析的对象之间的引用关系图谱。 Concurrent Remap:并发重映射这个阶段除了标记根对象直接引用的对象外,还会根据上个ZGC Cycle中生成的Forwarding Tables更新跨Heap Region的引用; Concurrent Mark:并发标记阶段,此阶段不会进行stw机制,会处于并发标记与根节点有关系并且可到达的对象。 同步检查点(Sync point):清理弱引用的数据信息 活动对象被移动到了一个新的Heap RegionB-region 中,之前旧对象所在的Heap RegionA-region 即可复用;如果 B-region 中对象之间的引用关系将会在这一阶段被更新; 在重定位过程中,新旧对象的映射关系(同一对象在不同Region中的映射关系)被记录在了Forwarding Tables中。 这允许我们在移动对象/整理内存阶段,在指向可回收/重用区域的指针确定之前回收/重用这部分内存,这有助于降低堆开销。这还意味着不需要实现单独的标记压缩算法来处理完整的GC。 这允许我们使用相对较少且简单的GC屏障。这有助于降低运行时开销。这还意味着在解释器和JIT编译器中更容易实现、优化和维护GC barrier代码。 局限性 ZGC的初始实验版本将不支持类卸载。默认情况下,classunload和ClassUnloadingWithConcurrentMark选项将被禁用。即便你启用也是不生效的。此外,ZGC最初不支持JVMCI(即Graal)。如果启用EnableJVMCI选项,将打印一条错误消息。 构建和使用 按照惯例,构建系统默认禁用JVM中的实验性特性。ZGC是一个实验性特性,因此不会出现在JDK构建中,除非在编译时使用configure选项: --with-jvm-features=zgc //显式地启用它。 (ZGC将出现在Oracle发布的所有Linux/x64 JDK版本中) JVM中的实验特性还需要在运行时显式地解锁。因此,要启用/使用ZGC,需要以下JVM选项: -XX:+ unlockexperimental alvmoptions -XX:+UseZGC 其他实现低延迟的方案 一种显而易见的方法是让G1的内存整理阶段可以并发执行。这种方案在很早的阶段就被放弃掉了。原因是现有的G1代码并不是为这种设计而编写的,因此要想在保证稳定的前提下实现这种特性非常困难。 一种方案是对现有的CMS收集器进行改进。然而这并不是一种好的选择,原因有很多,比如CMS不支持内存整理、代码复杂性和CMS实际上已经停止开发等。 Shenandoah Project正在探索通过使用Brooks Pointers来实现并发。 ZGC 总结 ZGC 在内存整理和引用更新上采取了不同的策略,给垃圾回收过程带来了巨大的性能提升。内存整理和引用更新都是并发的,也是交替进行的(其他的垃圾回收算法在更新引用时需要所有的线程到达safe-point)。但与此同时,我们也应该看到,并发带来的 GC 吞吐率的下降也是不可忽视的。 当响应时间比吞吐量占有更高的优先级时,ZGC 是个不错的选择。而对那些不能接受长时间暂停的应用程序来说,ZGC 是个理想的选择。而对于那些只是在后台进行密集计算的应用程序,G1 或者 Parallel 垃圾回收器可能具有更好的垃圾回收性能。 未完待续......

优秀的个人博客,低调大师

[雪峰磁针石博客]python tkinter图形工具样式作业

python测试开发项目实战-目录 python工具书籍下载-持续更新 使用tkinter绘制如下窗口 参考资料 本文最新版本地址 本文涉及的python测试开发库 谢谢点赞! 本文相关海量书籍下载 https://github.com/CoderDojoSV/beginner-python 代码 #!/usr/bin/env python3 # -*- coding: utf-8 -*- # 技术支持:https://www.jianshu.com/u/69f40328d4f0 # 技术支持 https://china-testing.github.io/ # https://github.com/china-testing/python-api-tesing/blob/master/practices/tk/tk4.py # 项目实战讨论QQ群630011153 144081101 # CreateDate: 2018-12-02 import tkinter as tk root = tk.Tk() root.configure(background='#4D4D4D') #top level styling # connecting to the external styling optionDB.txt root.option_readfile('optionDB.txt') #widget specific styling text = tk.Text( root, background='#101010', foreground="#D6D6D6", borderwidth=18, relief='sunken', width=17, height=5) text.insert( tk.END, "Style is knowing who you are,what you want to say, and not giving a damn." ) text.grid(row=0, column=0, columnspan=6, padx=5, pady=5) # all the below widgets derive their styling from optionDB.txt file tk.Button(root, text='*').grid(row=1, column=1) tk.Button(root, text='^').grid(row=1, column=2) tk.Button(root, text='#').grid(row=1, column=3) tk.Button(root, text='<').grid(row=2, column=1) tk.Button( root, text='OK', cursor='target').grid( row=2, column=2) #changing cursor style tk.Button(root, text='>').grid(row=2, column=3) tk.Button(root, text='+').grid(row=3, column=1) tk.Button(root, text='v').grid(row=3, column=2) tk.Button(root, text='-').grid(row=3, column=3) for i in range(10): tk.Button( root, text=str(i)).grid( column=3 if i % 3 == 0 else (1 if i % 3 == 1 else 2), row=4 if i <= 3 else (5 if i <= 6 else 6)) root.mainloop() 可以使用十六进制颜色代码为红色(r),绿色(g)和蓝色(b)的比例指定颜色。常用的表示是#rgb(4位),#rrggbb(8位)和#rrrgggbbb(12位)。 例如,#ff是白色,#000000是黑色,#f00是红色(R = 0xf,G = 0x0,B = 0x0),#00ff00为绿色(R = 0x00,G = 0xff,B = 0x00),#000000fff为蓝色(R = 0x000,G = 0x000,B = 0xfff)。 或者,Tkinter提供标准颜色名称的映射。有关预定义命名颜色的列表,请访问http://wiki.tcl.tk/37701或http://wiki.tcl.tk/16166。

优秀的个人博客,低调大师

[雪峰磁针石博客]python GUI作业:tkinter grid布局

python测试开发项目实战-目录 python工具书籍下载-持续更新 python 3.7极速入门教程 - 目录 要求 使用tkinter生成如下窗口: 参考资料 本文最新版本地址 本文涉及的python测试开发库 谢谢点赞! 本文相关海量书籍下载 python工具书籍下载-持续更新 python GUI工具书籍下载-持续更新 参考代码 #!/usr/bin/python3 # -*- coding: utf-8 -*- # 技术支持:https://www.jianshu.com/u/69f40328d4f0 # 技术支持 https://china-testing.github.io/ # https://github.com/china-testing/python-api-tesing/blob/master/practices/tk/tk2.py # 项目实战讨论QQ群630011153 144081101 # CreateDate: 2018-11-27 import tkinter as tk from tkinter import ttk from tkinter import scrolledtext from tkinter import Menu # Create instance win = tk.Tk() # Add a title win.title("Python GUI") tabControl = ttk.Notebook(win) # Create Tab Control tab1 = ttk.Frame(tabControl) # Create a tab tabControl.add(tab1, text='Tab 1') # Add the tab tab2 = ttk.Frame(tabControl) # Add a second tab tabControl.add(tab2, text='Tab 2') # Make second tab visible tabControl.pack(expand=1, fill="both") # Pack to make visible # LabelFrame using tab1 as the parent mighty = ttk.LabelFrame(tab1, text=' Mighty Python ') mighty.grid(column=0, row=0, padx=8, pady=4) # Modify adding a Label using mighty as the parent instead of win a_label = ttk.Label(mighty, text="Enter a name:") a_label.grid(column=0, row=0, sticky='W') # Modified Button Click Function def click_me(): action.configure(text='Hello ' + name.get() + ' ' + number_chosen.get()) # Adding a Textbox Entry widget name = tk.StringVar() name_entered = ttk.Entry(mighty, width=12, textvariable=name) name_entered.grid(column=0, row=1, sticky='W') # align left/West # Adding a Button action = ttk.Button(mighty, text="Click Me!", command=click_me) action.grid(column=2, row=1) # Creating three checkbuttons ttk.Label(mighty, text="Choose a number:").grid(column=1, row=0) number = tk.StringVar() number_chosen = ttk.Combobox(mighty, width=12, textvariable=number, state='readonly') number_chosen['values'] = (1, 2, 4, 42, 100) number_chosen.grid(column=1, row=1) number_chosen.current(0) chVarDis = tk.IntVar() check1 = tk.Checkbutton(mighty, text="Disabled", variable=chVarDis, state='disabled') check1.select() check1.grid(column=0, row=4, sticky=tk.W) chVarUn = tk.IntVar() check2 = tk.Checkbutton(mighty, text="UnChecked", variable=chVarUn) check2.deselect() check2.grid(column=1, row=4, sticky=tk.W) chVarEn = tk.IntVar() check3 = tk.Checkbutton(mighty, text="Enabled", variable=chVarEn) check3.deselect() check3.grid(column=2, row=4, sticky=tk.W) # GUI Callback function def checkCallback(*ignoredArgs): # only enable one checkbutton if chVarUn.get(): check3.configure(state='disabled') else: check3.configure(state='normal') if chVarEn.get(): check2.configure(state='disabled') else: check2.configure(state='normal') # trace the state of the two checkbuttons chVarUn.trace('w', lambda unused0, unused1, unused2 : checkCallback()) chVarEn.trace('w', lambda unused0, unused1, unused2 : checkCallback()) # Using a scrolled Text control scrol_w = 30 scrol_h = 3 scr = scrolledtext.ScrolledText(mighty, width=scrol_w, height=scrol_h, wrap=tk.WORD) scr.grid(column=0, row=5, sticky='WE', columnspan=3) # First, we change our Radiobutton global variables into a list colors = ["Blue", "Gold", "Red"] # We have also changed the callback function to be zero-based, using the list # instead of module-level global variables # Radiobutton Callback def radCall(): radSel=radVar.get() win.configure(background=colors[radSel]) # create three Radiobuttons using one variable radVar = tk.IntVar() # Next we are selecting a non-existing index value for radVar radVar.set(99) # Now we are creating all three Radiobutton widgets within one loop for col in range(3): curRad = tk.Radiobutton(mighty, text=colors[col], variable=radVar, value=col, command=radCall) curRad.grid(column=col, row=5, sticky=tk.W) # row=5 ... SURPRISE! # Create a container to hold labels buttons_frame = ttk.LabelFrame(mighty, text=' Labels in a Frame ') buttons_frame.grid(column=0, row=7) # Place labels into the container element ttk.Label(buttons_frame, text="Label1").grid(column=0, row=0, sticky=tk.W) ttk.Label(buttons_frame, text="Label2").grid(column=1, row=0, sticky=tk.W) ttk.Label(buttons_frame, text="Label3").grid(column=2, row=0, sticky=tk.W) # Exit GUI cleanly def _quit(): win.quit() win.destroy() exit() # Creating a Menu Bar menu_bar = Menu(win) win.config(menu=menu_bar) # Add menu items file_menu = Menu(menu_bar, tearoff=0) file_menu.add_command(label="New") file_menu.add_separator() file_menu.add_command(label="Exit", command=_quit) menu_bar.add_cascade(label="File", menu=file_menu) # Add another Menu to the Menu Bar and an item help_menu = Menu(menu_bar, tearoff=0) help_menu.add_command(label="About") menu_bar.add_cascade(label="Help", menu=help_menu) name_entered.focus() # Place cursor into name Entry #====================== # Start GUI #====================== win.mainloop()

优秀的个人博客,低调大师

[雪峰磁针石博客]性能测试工具nGrinder介绍

安装 以linux,这里采用的版本是centos 6 64bit,性能测试工具不建议在Windows上部署。 下载: https://github.com/naver/ngrinder/releases/ 选择最后面的war包。 服务器端启动: # java -XX:MaxPermSize=200m -jar ngrinder-controller-3.4.war --port 8058 这样nGrinder的管理页面就部署好,你可以简单的把ngrinder-controller的功能理解为性能测试展示和控制,后面会进行详细介绍。 打开网址: http://183.131.22.113:8058 默认用户名和密码都为admin 注意:这里的"Remember me"是短暂停留,页面关闭之后还是需要重新登陆的。 登录后点击右上角的admin,选择"下

优秀的个人博客,低调大师

[雪峰磁针石博客]python库介绍-multiprocessing:多进程

简介 进程是运行的程序,每个进程有自己的系统状态,包含了内存、打开文件列表、程序计数器(跟踪执行的指令)、存储函数本地调用变量的堆栈。 使用os或subprocess可以创建新进程,比如:os.fork(), subprocess.Popen()。子进程和父进程是相互独立执行的。 interprocess communication (IPC)进程间的通信: 最常见的形式是基于消息传递(message passing)。message是原始字节的缓存,通过I/O channel比如网络socket和管道,使用原语比如send() and recv()来发送接收消息。次常用的有内存映射区:memory-mapped regions,见mmap模块,实际上是共享内存。 线程有自己的控制流和执行堆栈,但是共享系统资源和数据。 并发的难点:同步和数据共享。解决的方法一般是使用互斥锁。 write_lock = Lock() ... # Critical section where writing occurs write_lock.acquire() f.write("Here's some data.\n") f.write("Here's more data.\n") ... write_lock.release() python的并发程序设计 多数系统上,Python支持消息传递和基于线程的并发程序设计。global interpreter lock (the GIL)机制实际每个时间单元只允许单个线程执行,哪怕有多个CPU。如果瓶颈在I/O,使用多线程效果不错;如果在cpu,效果则会更差。还不如使用子进程和消息传递。线程数一多经常出现以下怪异的问题,比如100个线程工作良好,1000个线程就可能出问题了,这种情况一般需要使用异步事件处理系统,比如中央事件循环可能使用select模块监控I/O资源和分发异步到大量的I/O 处理器。asyncore和流行的第三方的Twisted (http://twistedmatrix/com)可以实现这点。 消息传递在python使用很广,甚至在线程中。它难于出错,减少了锁和同步原语的使用。可以扩展至网络和分布式系统。Python的高级特性比如协程序(coroutines)也使用消息传递抽象。 multiprocessing支持子进程、通信和共享数据、执行不同形式的同步。 multiprocessing Process类 这个类表示子进程中运行的任务:Process([group [, target [, name [, args [, kwargs]]]]]),构造函数中必须使用关键字参数,target表示可调用对象,args表示调用对象的位置参数元组。kwargs表示调用对象的字典。Name为别名。Group实质上不使用。 方法有:is_alive()、.join([timeout])、run()、start()、terminate()。 属性有:authkey、daemon(要通过start()设置)、exitcode(进程在运行时为None、如果为–N,表示被信号N结束)、name、pid。 Process类中,注意daemon是父进程终止后自动终止,且自己不能产生新进程,必须在start()之前设置。 创建函数并将其作为单个进程。 import multiprocessing import time def clock(interval): for i in range(3): print("The time is {0}".format(time.ctime())) time.sleep(interval) if __name__ == '__main__': p = multiprocessing.Process(target=clock, args=(2,)) p.start() 将进程定义为类: import multiprocessing import time class ClockProcess(multiprocessing.Process): def __init__(self, interval): multiprocessing.Process.__init__(self) self.interval = interval def run(self): for i in range(3): print("The time is {0}".format(time.ctime())) time.sleep(self.interval) if __name__ == '__main__': p = ClockProcess(2) p.start() 注意,要在命令行才能执行,用IDE是不行的。 进程通信 multiprocessing支持管道和队列,都是用消息传递来实现的,队列接口和线程中的队列类似。 Queue([maxsize]):默认不限制大小,队列实质是用管道和锁来实现的。支持线程会给底层管道传送数据。 方法有:cancel_join_thread()、close()、empty()、full()、get([block [, timeout]])、get_nowait()(等同于get(False))、join_thread()、put(item [, block [, timeout]])、put_nowait(item)(等同于put(item, False))、qsize()、JoinableQueue([maxsize])、task_done()、join() 下例使用队列进行通信: JoinableQueue创建连接的进程队列。队列和普通队列基本一样,不过消费者在处理完毕之后可以通知生产者(q.task_done())。使用共享信号和条件变量实现。join()由生产者使用,等待所有成员都收到task_done。 import multiprocessing def consumer(input_q): while True: item = input_q.get() print(item) input_q.task_done() def producer(sequence, output_q): for item in sequence: output_q.put(item) if __name__ == '__main__': q = multiprocessing.JoinableQueue() cons_p = multiprocessing.Process(target=consumer, args=(q,)) cons_p.daemon = True cons_p.start() sequence = [1, 2, 3, 4] producer(sequence, q) q.join() 这里控制多进程的关键在于队列get()之后,使用task_done()指示该元素处理完毕;进程启动之前设置了daemon为True;对队列使用join()。 这种方法可以启动多个进程,如下: process = [] key_list = multiprocessing.JoinableQueue() # Launch the consumer process for i in range(10): t = multiprocessing.Process(target=consumer,args=(key_list,lock)) t.daemon=True process.append(t) for i in range(10): process[i].start() producer( key_list ) key_list.join() 下面有个应用实例: https://bitbucket.org/china-testing/small_python_daily_tools/src/87d81739633482abdd3a2d0d11f62f6edd989555/db/mysql/check_transfer.py?at=default&fileviewer=file-view-default 在某些程序中,生产者需要告知消费者没有更多项目了,消费者可以关闭了。这时需要使用哨兵(sentinel)。 #!/usr/bin/env python # -*- coding: utf-8 -*- # multiprocessing_sentinel.py # Author Rongzhong Xu 2016-08-11 wechat: pythontesting """ multiprocessing sentinel demo, Tesed in python2.7/3.5/2.6 """ import multiprocessing def consumer(input_q): while True: item = input_q.get() if item is None: break # Process item print(item) # Replace with useful work # Shutdown print("Consumer done") def producer(sequence, output_q): for item in sequence: # Put the item on the queue output_q.put(item) if __name__ == '__main__': q = multiprocessing.Queue() # Launch the consumer process cons_p = multiprocessing.Process(target=consumer, args=(q,)) cons_p.start() # Produce items sequence = [1, 2, 3, 4] producer(sequence, q) # Signal completion by putting the sentinel on the queue q.put(None) # Wait for the consumer process to shutdown cons_p.join() 注意:每个消费者都需要一个:sentinel,可以使用for语句来实现 for i in range(10): q.put(None) 实际使用中不局限于使用None,使用其他特殊符号等也是可以的。上面程序从表面看比使用JoinableQueue要复杂,实现的效果又是一样的。实际上这种场景应用更广泛,在consumer比较耗时的情况下,JoinableQueue如果锁住整个函数则互相等待的时间太长,如果不锁,后面几次执行可能丢失数据。 管道 使用管道:Pipe([duplex]),返回值:元组(conn1, conn2)。conn1和conn2为Connection对象,代表管道的末端。管道默认是双向的,如果设置duplex为False,conn1只能接收,conn2只能发送。 Connection对象的方法和属性如下: close()、fileno()、poll([timeout])、recv()、recv_bytes([maxlength])、recv_bytes_into(buffer [, offset])、send(obj)、send_bytes(buffer [, offset [, size]]) 下面例子实现和之前类似的功能: def consumer(pipe): output_p, input_p = pipe input_p.close() # Close the input end of the pipe while True: try: item = output_p.recv() except EOFError: break # Process item print(item) # Replace with useful work # Shutdown print("Consumer done") # Produce items and put on a queue. sequence is an # iterable representing items to be processed. def producer(sequence, input_p): for item in sequence: # Put the item on the queue input_p.send(item) if __name__ == '__main__': (output_p, input_p) = multiprocessing.Pipe() # Launch the consumer process cons_p = multiprocessing.Process( target=consumer, args=((output_p, input_p),)) cons_p.start() # Close the output pipe in the producer output_p.close() # Produce items sequence = [1, 2, 3, 4] producer(sequence, input_p) # Signal completion by closing the input pipe input_p.close() # Wait for the consumer process to shutdown cons_p.join() 管道还可以用于双向通信,比如下例的C/S模式: import multiprocessing # A server process def adder(pipe): server_p, client_p = pipe client_p.close() while True: try: x, y = server_p.recv() except EOFError: break result = x + y server_p.send(result) # Shutdown print("Server done") if __name__ == '__main__': (server_p, client_p) = multiprocessing.Pipe() # Launch the server process adder_p = multiprocessing.Process( target=adder, args=((server_p, client_p),)) adder_p.start() # Close the server pipe in the client server_p.close() # Make some requests on the server client_p.send((3, 4)) print(client_p.recv()) client_p.send(('Hello', 'World')) print(client_p.recv()) # Done. Close the pipe client_p.close() # Wait for the consumer process to shutdown adder_p.join() send()和recv()使用pickle序列化对象。更高级的程序需要使用远程过程调用,需要使用到进程池。 进程池 Pool类在简单的情况下可用于管理固定数量的消费者。进程池的功能和列表解析及函数式编程中的map-reduce类似。 import multiprocessing import time def do_calculation(data): return data * 2 def start_process(): print('Starting {0}'.format(multiprocessing.current_process().name)) if __name__ == '__main__': # convert range to list for python3 inputs = list(range(100)) time1 = time.time() builtin_outputs = map(do_calculation, inputs) # convert to list for python3 print('Built-in: {0}'.format(list(builtin_outputs))) time2 = time.time() print(time2 - time1) pool_size = multiprocessing.cpu_count() * 2 pool = multiprocessing.Pool(processes=pool_size, initializer=start_process, ) pool_outputs = pool.map(do_calculation, inputs) pool.close() # no more tasks pool.join() # wrap up current tasks time3 = time.time() print('Pool : {0}'.format(pool_outputs)) print(time3 - time2) 执行结果: $ python3 multiprocessing_pool.py Built-in: [0, 2, 4, 6, 8, 10, 12, 14, 16, 18, 20, 22, 24, 26, 28, 30, 32, 34, 36, 38, 40, 42, 44, 46, 48, 50, 52, 54, 56, 58, 60, 62, 64, 66, 68, 70, 72, 74, 76, 78, 80, 82, 84, 86, 88, 90, 92, 94, 96, 98, 100, 102, 104, 106, 108, 110, 112, 114, 116, 118, 120, 122, 124, 126, 128, 130, 132, 134, 136, 138, 140, 142, 144, 146, 148, 150, 152, 154, 156, 158, 160, 162, 164, 166, 168, 170, 172, 174, 176, 178, 180, 182, 184, 186, 188, 190, 192, 194, 196, 198] 3.790855407714844e-05 Starting ForkPoolWorker-1 Starting ForkPoolWorker-2 Starting ForkPoolWorker-3 Starting ForkPoolWorker-4 Starting ForkPoolWorker-5 Starting ForkPoolWorker-6 Starting ForkPoolWorker-7 Starting ForkPoolWorker-8 Starting ForkPoolWorker-9 Starting ForkPoolWorker-10 Starting ForkPoolWorker-11 Starting ForkPoolWorker-12 Starting ForkPoolWorker-13 Starting ForkPoolWorker-14 Starting ForkPoolWorker-15 Starting ForkPoolWorker-16 Pool : [0, 2, 4, 6, 8, 10, 12, 14, 16, 18, 20, 22, 24, 26, 28, 30, 32, 34, 36, 38, 40, 42, 44, 46, 48, 50, 52, 54, 56, 58, 60, 62, 64, 66, 68, 70, 72, 74, 76, 78, 80, 82, 84, 86, 88, 90, 92, 94, 96, 98, 100, 102, 104, 106, 108, 110, 112, 114, 116, 118, 120, 122, 124, 126, 128, 130, 132, 134, 136, 138, 140, 142, 144, 146, 148, 150, 152, 154, 156, 158, 160, 162, 164, 166, 168, 170, 172, 174, 176, 178, 180, 182, 184, 186, 188, 190, 192, 194, 196, 198] 0.2203056812286377 上面例子先计算map的时间,然后用进程池的map,计算出时间。在列表数比较少的情况下,多进程的执行时间更短。列表数比较多的情况下,多进程的执行时间更长,可见python内置的map是效率比较高的。 如果消费者函数有内存泄露,可以在执行任务之后重启,设定maxtasksperchild参数即可。 import time def do_calculation(data): return data * 2 def start_process(): print('Starting {0}'.format(multiprocessing.current_process().name)) if __name__ == '__main__': # convert range to list for python3 inputs = list(range(100)) time1 = time.time() builtin_outputs = map(do_calculation, inputs) # convert to list for python3 print('Built-in: {0}'.format(list(builtin_outputs))) time2 = time.time() print(time2 - time1) pool_size = multiprocessing.cpu_count() * 2 pool = multiprocessing.Pool(processes=pool_size, initializer=start_process, maxtasksperchild=3, ) pool_outputs = pool.map(do_calculation, inputs) pool.close() # no more tasks pool.join() # wrap up current tasks time3 = time.time() print('Pool : {0}'.format(pool_outputs)) print(time3 - time2) 执行结果: $ python3 multiprocessing_pool2.py Built-in: [0, 2, 4, 6, 8, 10, 12, 14, 16, 18, 20, 22, 24, 26, 28, 30, 32, 34, 36, 38, 40, 42, 44, 46, 48, 50, 52, 54, 56, 58, 60, 62, 64, 66, 68, 70, 72, 74, 76, 78, 80, 82, 84, 86, 88, 90, 92, 94, 96, 98, 100, 102, 104, 106, 108, 110, 112, 114, 116, 118, 120, 122, 124, 126, 128, 130, 132, 134, 136, 138, 140, 142, 144, 146, 148, 150, 152, 154, 156, 158, 160, 162, 164, 166, 168, 170, 172, 174, 176, 178, 180, 182, 184, 186, 188, 190, 192, 194, 196, 198] 3.600120544433594e-05 Starting ForkPoolWorker-1 Starting ForkPoolWorker-3 Starting ForkPoolWorker-2 Starting ForkPoolWorker-4 Starting ForkPoolWorker-5 Starting ForkPoolWorker-6 Starting ForkPoolWorker-7 Starting ForkPoolWorker-8 Starting ForkPoolWorker-9 Starting ForkPoolWorker-10 Starting ForkPoolWorker-11 Starting ForkPoolWorker-12 Starting ForkPoolWorker-13 Starting ForkPoolWorker-14 Starting ForkPoolWorker-15 Starting ForkPoolWorker-16 Starting ForkPoolWorker-17 Starting ForkPoolWorker-18 Starting ForkPoolWorker-19 Starting ForkPoolWorker-20 Starting ForkPoolWorker-21 Starting ForkPoolWorker-22 Starting ForkPoolWorker-23 Starting ForkPoolWorker-24 Starting ForkPoolWorker-25 Starting ForkPoolWorker-26 Starting ForkPoolWorker-27 Starting ForkPoolWorker-28 Starting ForkPoolWorker-29 Starting ForkPoolWorker-30 Starting ForkPoolWorker-31 Starting ForkPoolWorker-32 Pool : [0, 2, 4, 6, 8, 10, 12, 14, 16, 18, 20, 22, 24, 26, 28, 30, 32, 34, 36, 38, 40, 42, 44, 46, 48, 50, 52, 54, 56, 58, 60, 62, 64, 66, 68, 70, 72, 74, 76, 78, 80, 82, 84, 86, 88, 90, 92, 94, 96, 98, 100, 102, 104, 106, 108, 110, 112, 114, 116, 118, 120, 122, 124, 126, 128, 130, 132, 134, 136, 138, 140, 142, 144, 146, 148, 150, 152, 154, 156, 158, 160, 162, 164, 166, 168, 170, 172, 174, 176, 178, 180, 182, 184, 186, 188, 190, 192, 194, 196, 198] 0.23842501640319824 从结果看,进程数有所增加。(注意,进程数似乎比预期的要少) Pool([numprocess [,initializer [, initargs]]]) numprocess的默认值是cpu_count()。方法有:apply(func [, args [, kwargs]]),apply_async(func [, args [, kwargs [, callback]]]),close(),join(),imap(func, iterable [, chunksize]),imap_unordered(func, iterable [, chunksize]]),map(func, iterable [, chunksize]),map_async(func, iterable [, chunksize [, callback]]),terminate(). 返回结果AsyncResult的方法:get([timeout])、ready()、sucessful()、wait([timeout])、wait([timeout]) 以下代码生成指定目录的文件名和SHA512对应表的字典。 import multiprocessing import hashlib import binascii # Some parameters you can tweak BUFSIZE = 8192 # Read buffer size POOLSIZE = 2 # Number of workers def compute_digest(filename): try: f = open(filename, "rb") except IOError: return None digest = hashlib.sha512() while True: chunk = f.read(BUFSIZE) if not chunk: break digest.update(chunk) f.close() return filename, digest.digest() def build_digest_map(topdir): digest_pool = multiprocessing.Pool(POOLSIZE) allfiles = (os.path.join(path, name) for path, dirs, files in os.walk(topdir) for name in files) digest_map = dict(digest_pool.imap_unordered(compute_digest, allfiles, 20)) digest_pool.close() return digest_map # Try it out. Change the directory name as desired. if __name__ == '__main__': digest_map = build_digest_map("/home/andrew/data/code/python/\ python-chinese-library/libraries/multiprocessing") print(len(digest_map)) for key in digest_map.keys(): print("{0}: {1}".format(key, binascii.hexlify(digest_map[key]))) 共享数据和同步 共享内存通过mmap实现。共享内存中创建的是ctypes对象,不需要管道中的序列化。 Value(typecode, arg1, ... argN, lock),RawValue(typecode, arg1, ..., argN),Array(typecode, initializer, lock),RawArray(typecode, initializer) 原语有: Lock,Rlock,Semaphore,BoundedSemaphore,Event,Condition. import multiprocessing class FloatChannel(object): def __init__(self, maxsize): self.buffer = multiprocessing.RawArray('d', maxsize) self.buffer_len = multiprocessing.Value('i') self.empty = multiprocessing.Semaphore(1) self.full = multiprocessing.Semaphore(0) def send(self, values): self.empty.acquire() # Only proceed if buffer empty nitems = len(values) self.buffer_len = nitems # Set the buffer size self.buffer[:nitems] = values # Copy values into the buffer self.full.release() # Signal that buffer is full def recv(self): self.full.acquire() # Only proceed if buffer full values = self.buffer[:self.buffer_len.value] # Copy values self.empty.release() # Signal that buffer is empty return values # Performance test. Receive a bunch of messages def consume_test(count, ch): for i in range(count): values = ch.recv() # Performance test. Send a bunch of messages def produce_test(count, values, ch): for i in range(count): ch.send(values) if __name__ == '__main__': ch = FloatChannel(100000) p = multiprocessing.Process(target=consume_test, args=(1000, ch)) p.start() values = [float(x) for x in range(100000)] produce_test(1000, values, ch) print("Done") p.join() 参考资料 紧张整理更新中,讨论 钉钉免费群21745728 qq群144081101 567351477 本文涉及的python测试开发库 谢谢点赞! 本文最新版本地址 本文相关书籍下载 pymotw multiprocessing参考

资源下载

更多资源
腾讯云软件源

腾讯云软件源

为解决软件依赖安装时官方源访问速度慢的问题,腾讯云为一些软件搭建了缓存服务。您可以通过使用腾讯云软件源站来提升依赖包的安装速度。为了方便用户自由搭建服务架构,目前腾讯云软件源站支持公网访问和内网访问。

Nacos

Nacos

Nacos /nɑ:kəʊs/ 是 Dynamic Naming and Configuration Service 的首字母简称,一个易于构建 AI Agent 应用的动态服务发现、配置管理和AI智能体管理平台。Nacos 致力于帮助您发现、配置和管理微服务及AI智能体应用。Nacos 提供了一组简单易用的特性集,帮助您快速实现动态服务发现、服务配置、服务元数据、流量管理。Nacos 帮助您更敏捷和容易地构建、交付和管理微服务平台。

Rocky Linux

Rocky Linux

Rocky Linux(中文名:洛基)是由Gregory Kurtzer于2020年12月发起的企业级Linux发行版,作为CentOS稳定版停止维护后与RHEL(Red Hat Enterprise Linux)完全兼容的开源替代方案,由社区拥有并管理,支持x86_64、aarch64等架构。其通过重新编译RHEL源代码提供长期稳定性,采用模块化包装和SELinux安全架构,默认包含GNOME桌面环境及XFS文件系统,支持十年生命周期更新。

WebStorm

WebStorm

WebStorm 是jetbrains公司旗下一款JavaScript 开发工具。目前已经被广大中国JS开发者誉为“Web前端开发神器”、“最强大的HTML5编辑器”、“最智能的JavaScript IDE”等。与IntelliJ IDEA同源,继承了IntelliJ IDEA强大的JS部分的功能。

用户登录
用户注册