Python常用的工具代码


import logging
import requests
import re
from lxml import etree
import sql_crud

headers = {
    'user-agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) '
                  'Chrome/96.0.4664.45 Safari/537.36'
}
LOG_FORMAT = "%(levelname)s - %(asctime)s - %(message)s"
logging.basicConfig(level=logging.INFO, format=LOG_FORMAT)


class ZjzwApiMonitor:
    def __init__(self, front_page):
        self.front_page = front_page

    """处理html,将取得的url去重"""

    def __get_html_url(self, url):

        if not url.find('.exe') and not url.find('.pdf') and not url.find('.mp4') and not url.find(
                '.zip') and not url.find('.doc'):
            return None
        logging.info("正在get:%s" % url)
        # html = requests.get(url, headers=headers)
        url_list = set()
        match_zjzw = re.compile(r'^w.*\.zj\..*|^h.*\.zj\..*|^w.*zjzw.*|^h.*zjzw.*')
        try:
            html = requests.get(url, headers=headers, verify=False, timeout=3)
            etree_html = etree.HTML(html.text, etree.HTMLParser())
            href = etree_html.xpath('//*/@href')
            for line in href:
                match_list = match_zjzw.findall(line)
                if match_list:
                    url_list.add(match_list[0])
        except AttributeError as ex:
            logging.error(repr(ex))
        except requests.exceptions.ConnectTimeout as ex:
            logging.error(repr(ex))
        except requests.exceptions.InvalidURL as ex:
            logging.error(repr(ex))
        except requests.exceptions.ReadTimeout as ex:
            logging.error(repr(ex))
        except requests.exceptions.ConnectionError as ex:
            logging.error(repr(ex))
        return url_list

    """获取数据库的url,请求这些url再次写入数据库"""

    def control(self):
        url_set, found_set = sql_crud.select()
        # 将首页url插入数据库
        if self.front_page not in found_set:
            url_list = self.__get_html_url(self.front_page)
            if not url_list:
                sql_crud.insert(url_list, self.front_page)
        # 取两列差集,判断该url有没有get过
        found = url_set - found_set
        # 从html中取得的url插入数据库
        for url in found:
            html_url = self.__get_html_url(url)
            # 用数据库中的url和页面的url比对,返回差集
            unique_url = html_url - url_set
            sql_crud.insert(unique_url, url)

    def __get_status_code(self):
        # 无法get状态码的url需要异常捕获
        # 获取请求时间
        pass


if __name__ == '__main__':
    domain = ZjzwApiMonitor("https://www.zj.gov.cn/")
    domain.control()

import sys
import pymysql
import configparser
import logging

CONFIG = configparser.ConfigParser()
CONFIG.read('db.cfg')
LOG_FORMAT = "%(levelname)s - %(asctime)s -  %(message)s"
logging.basicConfig(level=logging.INFO, format=LOG_FORMAT)
HOST = CONFIG.get('mysql', 'host')
USER = CONFIG.get('mysql', 'user')
PASSWORD = CONFIG.get('mysql', 'password')
DATABASE = CONFIG.get('mysql', 'database')

"""数据插入"""


def insert(url_list, found):
    """连接mysql数据库"""
    conn = pymysql.connect(host=HOST, user=USER, password=PASSWORD, database=DATABASE)
    logging.debug("插入 %s 去重后的URL" % found)
    sql = 'INSERT INTO total_url(`ID`,`url`,`date`,`found`) VALUES(null,%s,NOW(),%s)'
    data_list = []
    try:
        cursor = conn.cursor()
        for url in url_list:
            data_list.append((url, found))
        cursor.executemany(sql, data_list)
        conn.commit()
    except pymysql.InterfaceError as ex:
        logging.error("%s 数据库连接断开,尝试重连..." % repr(ex))
        conn.ping(reconnect=True)
    except pymysql.IntegrityError as ex:
        logging.error("重复插入! %s" % repr(ex))
    except Exception as ex:
        logging.error(repr(ex))
    conn.close()


"""返回列"""


def select():
    """连接mysql数据库"""
    conn = pymysql.connect(host=HOST, user=USER, password=PASSWORD, database=DATABASE)
    cursor = conn.cursor()
    sql = "SELECT `url`,`found` FROM total_url"
    difference_set_sql = "SELECT count(url) FROM total_url WHERE url NOT IN (SELECT DISTINCT " \
                         "`found` FROM total_url WHERE `found` IS NOT NULL);"

    url_set = set()
    found_set = set()
    try:
        cursor.execute(sql)
        results = cursor.fetchall()
        for u in results:
            url_set.add(u[0])
            found_set.add(u[1])
    except Exception as ex:
        logging.error(repr(ex))
    cursor.close()
    conn.close()
    return url_set, found_set

# def test():
#     cursor = conn.cursor()
#     n = 'NULL'
#     sql = "insert into total_url(`ID`,`url`,`date`,`found`) VALUES(null,'http',NOW(),%s)"
#     cursor.execute(sql, n)
#     conn.commit()
#     cursor.close()
#     conn.close()

调用mysql

#!/usr/bin/python3
 
import pymysql
 
# 打开数据库连接
db = pymysql.connect(host='localhost',
                     user='testuser',
                     password='test123',
                     database='TESTDB')
 
# 使用 cursor() 方法创建一个游标对象 cursor
cursor = db.cursor()
 
# 使用 execute()  方法执行 SQL 查询 
cursor.execute("SELECT VERSION()")
 
# 使用 fetchone() 方法获取单条数据.
data = cursor.fetchone()
 
print ("Database version : %s " % data)
 
# 关闭数据库连接
db.close()

image-20211126145942530

创建数据库表

import pymysql
 
# 打开数据库连接
db = pymysql.connect(host='localhost',
                     user='testuser',
                     password='test123',
                     database='TESTDB')
 
# 使用 cursor() 方法创建一个游标对象 cursor
cursor = db.cursor()
 
# 使用 execute() 方法执行 SQL,如果表存在则删除
cursor.execute("DROP TABLE IF EXISTS EMPLOYEE")
 
# 使用预处理语句创建表
sql = """CREATE TABLE EMPLOYEE (
         FIRST_NAME  CHAR(20) NOT NULL,
         LAST_NAME  CHAR(20),
         AGE INT,  
         SEX CHAR(1),
         INCOME FLOAT )"""
 
cursor.execute(sql)
 
# 关闭数据库连接
db.close()

数据库插入操作

import pymysql
 
# 打开数据库连接
db = pymysql.connect(host='localhost',
                     user='testuser',
                     password='test123',
                     database='TESTDB')
 
# 使用cursor()方法获取操作游标 
cursor = db.cursor()
 
# SQL 插入语句
sql = "insert into `bird_record`(`id`,`name`,`ph`) values(%s,%s,%s)"

val = (('li', 'si', 16, 'F', 1000),
       ('Bruse', 'Jerry', 30, 'F', 3000),
       ('Lee', 'Tomcat', 40, 'M', 4000),
       ('zhang', 'san', 18, 'M', 1500))


try:
   # 执行sql语句
   cursor.execute(sql, (参数1,参数2···)
   #批量插入
   cursor.executemany(sql,val)

   # 提交到数据库执行
   db.commit()
except:
   # 如果发生错误则回滚
   db.rollback()
    
# 关闭数据库连接
db.close()

以上例子也可以写成如下形式:

with connection.cursor() as cursor:
    cursor.executemany("insert into test(prop, val) values (%s, %s)", vals,name)
    connection.commit()

数据库查询操作

Python查询Mysql使用 fetchone() 方法获取单条数据, 使用fetchall() 方法获取多条数据。

  • fetchone(): 该方法获取下一个查询结果集。结果集是一个对象
  • fetchall(): 接收全部的返回结果行.
  • rowcount: 这是一个只读属性,并返回执行execute()方法后影响的行数。
import pymysql
 
# 打开数据库连接
db = pymysql.connect(host='localhost',
                     user='testuser',
                     password='test123',
                     database='TESTDB')
 
# 使用cursor()方法获取操作游标 
cursor = db.cursor()
 
# SQL 查询语句
sql = "SELECT * FROM EMPLOYEE \
       WHERE INCOME > %s" % (1000)
try:
   # 执行SQL语句
   cursor.execute(sql)
   # 获取所有记录列表
   results = cursor.fetchall()
   for row in results:
      fname = row[0]
      lname = row[1]
      age = row[2]
      sex = row[3]
      income = row[4]
       # 打印结果
      print ("fname=%s,lname=%s,age=%s,sex=%s,income=%s" % \
             (fname, lname, age, sex, income ))
except:
   print ("Error: unable to fetch data")
 
# 关闭数据库连接
db.close()

数据库更新操作

import pymysql
 
# 打开数据库连接
db = pymysql.connect(host='localhost',
                     user='testuser',
                     password='test123',
                     database='TESTDB')
 
# 使用cursor()方法获取操作游标 
cursor = db.cursor()
 
# SQL 更新语句
sql = "UPDATE EMPLOYEE SET AGE = AGE + 1 WHERE SEX = '%c'" % ('M')
try:
   # 执行SQL语句
   cursor.execute(sql)
   # 提交到数据库执行
   db.commit()
except:
   # 发生错误时回滚
   db.rollback()
 
# 关闭数据库连接
db.close()

数据库删除操作

import pymysql
 
# 打开数据库连接
db = pymysql.connect(host='localhost',
                     user='testuser',
                     password='test123',
                     database='TESTDB')
 
# 使用cursor()方法获取操作游标 
cursor = db.cursor()
 
# SQL 删除语句
sql = "DELETE FROM EMPLOYEE WHERE AGE > %s" % (20)
try:
   # 执行SQL语句
   cursor.execute(sql)
   # 提交修改
   db.commit()
except:
   # 发生错误时回滚
   db.rollback()
 
# 关闭连接
db.close()

执行事务

事务机制可以确保数据一致性。

事务应该具有4个属性:原子性、一致性、隔离性、持久性。这四个属性通常称为ACID特性。

  • 原子性(atomicity)。一个事务是一个不可分割的工作单位,事务中包括的诸操作要么都做,要么都不做。
  • 一致性(consistency)。事务必须是使数据库从一个一致性状态变到另一个一致性状态。一致性与原子性是密切相关的。
  • 隔离性(isolation)。一个事务的执行不能被其他事务干扰。即一个事务内部的操作及使用的数据对并发的其他事务是隔离的,并发执行的各个事务之间不能互相干扰。
  • 持久性(durability)。持续性也称永久性(permanence),指一个事务一旦提交,它对数据库中数据的改变就应该是永久性的。接下来的其他操作或故障不应该对其有任何影响。
# SQL删除记录语句
sql = "DELETE FROM EMPLOYEE WHERE AGE > %s" % (20)
try:
   # 执行SQL语句
   cursor.execute(sql)
   # 向数据库提交
   db.commit()
except:
   # 发生错误时回滚
   db.rollback()

对于支持事务的数据库, 在Python数据库编程中,当游标建立之时,就自动开始了一个隐形的数据库事务。

commit()方法游标的所有更新操作,rollback()方法回滚当前游标的所有操作。每一个方法都开始了一个新的事务。

错误处理

DB API中定义了一些数据库操作的错误及异常,下表列出了这些错误和异常:

异常 描述
Warning 当有严重警告时触发,例如插入数据是被截断等等。必须是 StandardError 的子类。
Error 警告以外所有其他错误类。必须是 StandardError 的子类。
InterfaceError 当有数据库接口模块本身的错误(而不是数据库的错误)发生时触发。 必须是Error的子类。
DatabaseError 和数据库有关的错误发生时触发。 必须是Error的子类。
DataError 当有数据处理时的错误发生时触发,例如:除零错误,数据超范围等等。 必须是DatabaseError的子类。
OperationalError 指非用户控制的,而是操作数据库时发生的错误。例如:连接意外断开、 数据库名未找到、事务处理失败、内存分配错误等等操作数据库是发生的错误。 必须是DatabaseError的子类。
IntegrityError 完整性相关的错误,例如外键检查失败等。必须是DatabaseError子类。
InternalError 数据库的内部错误,例如游标(cursor)失效了、事务同步失败等等。 必须是DatabaseError子类。
ProgrammingError 程序错误,例如数据表(table)没找到或已存在、SQL语句语法错误、 参数数量错误等等。必须是DatabaseError的子类。
NotSupportedError 不支持错误,指使用了数据库不支持的函数或API等。例如在连接对象上 使用.rollback()函数,然而数据库并不支持事务或者事务已关闭。 必须是DatabaseError的子类。

远程Linux

import paramiko

client = paramiko.SSHClient()
# client.load_host_keys()
client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
client.connect(hostname='192.168.111.141', port=22, username='root', password='123')
stdin, stdout, stderr = client.exec_command('ls /mnt')
print("stdout:", stdout.readline())
# print("stderr:", stderr.read())
# print("stdin:",stdin.read())

生成测试数据

mimesis

mimesis 是一个高性能的伪数据生成器,目前支持 33 种不同的语言环境。通过该库,我们可以生成各种测试数据、假的 API 接口、任意结构的 JSON 和 XML 数据以及隐藏生产环境的数据。

pip install mimesis 

安装好之后我们就可以直接使用了。

from mimesis import Person

person = Person('zh')
print(f'name: {person.surname() + "" + person.name()}')
print(f'sex: {person.sex()}')
print(f'academic degree: {person.academic_degree()}')

## 输出结果
name: 田曜岩
sex: 男性
academic degree: 研究生

在上面的程序中,我们创建了一个使用中文环境的 Person 对象,接着输出该用户的姓名,性别以及学历。

下面我们看看 Person 对象里面都有啥假数据。

print('\n'.join(('%s:%s' % item for item in person._data.items()))) 

结果如下所示:

img

除了姓名,性别这些基本信息之外还有学历、性取向、大学以及信仰等信息。

另外,除了 Person 之外,mimesis 库还提供了 Address、Food、Datetime 等方面的数据。

address = Address("zh")
print(f'continent: {address.continent()}')
print(f'province: {address.province()}')
print(f'city: {address.city()}')
print(f'street name: {address.street_name()}')

## 输出结果
province: 安徽省
city: 湛江市

image-20211228165943128

除了省份,城市之外还有大陆、国家、州、街区等信息。

food = Food("zh")
print(f'dish: {food.dish()}')
print(f'drink: {food.drink()}')

## 输出结果
dish: 东坡肉
drink: 红茶

image-20211228170126076

除了鱼类和饮料之外还有水果、香料和蔬菜。

其实 mimesis 库的强大不止于此,甚至我们可以使用该库来返回特定格式的数据。这就要借助 mimesis.schema 来实现了。

比如,我们要返回如下格式的 JSON 数据,那么就可以这么写:

_ = Field('zh')
schema = Schema(schema=lambda: {
    'id': _('uuid'),
    'name': _('person.name'),
    'version': _('version', pre_release=True),
    'timestamp': _('timestamp', posix=False),
    'owner': {
        'email': _('person.email', domains=['test.com'], key=str.lower),
        'token': _('token_hex'),
        'creator': _('full_name', gender=Gender.FEMALE)
    },
    'address': {
        'country': _('address.country'),
        'province': _('address.province'),
        'city': _('address.city')
    }
})

# 生成数据
data = schema.create(iterations=2)

我们借助 Flask 快速实现一个接口:

@app.route('/apps', methods=('GET',))
def apps_view():
    count = request.args.get('count', default=1, type=int)
    data = schema.create(iterations=count)
    return jsonify(data)

调用该接口获得如下数据:

image-20211228170141944

faker

faker 同样是一个优秀的生成假数据的 Python 库,支持多种语言环境,我们可以使用 pip 进行安装。

pip install faker 

试着获取一下姓名,地址,日期等假数据

from faker import Faker

faker = Faker(locale='zh_CN')
print(f'name: {faker.name()}')
print(f'address: {faker.address()}')
print(f'date: {faker.date()}')

## 输出结果
name: 刘晶
address: 香港特别行政区平县华龙深圳路l座 580988
date: 2013-08-26
image-20211228170158734

从上图可以看出,faker 目前支持 22 个不同种类的假数据。

如果以上类别不能满足需求,那么 faker 同样支持自定义扩展。

from faker.providers import BaseProvider

class MyProvider(BaseProvider):
    def foo(self):
        return 'bar'
faker.add_provider(MyProvider)
print(f'foo: {faker.foo()}')

## 输出结果
foo: bar

最方便的是我们可以直接在命令行调用 faker,这对于某些场景简直不要太方便,比如我们就需要一条用户信息的 JSON 数据,那么就不需要再去写一个 Python 脚本了,直接在命令行调用 faker 命令即可生成假数据。

$ faker address
香港特别行政区六安市海港哈尔滨街z座 561730

$ faker -r=3 address
湖南省海口市清浦王街h座 140394
海南省银川市孝南武汉街q座 623233
青海省建华市萧山李街w座 207439

$ faker -r=3 profile name,address,birthdate
{'name': '邵丽', 'address': '四川省想市上街吴街c座 399962', 'birthdate': datetime.date(1979, 8, 28)};
{'name': '张秀华', 'address': '江西省亮市徐汇程街p座 527720', 'birthdate': datetime.date(1907, 6, 27)};
{'name': '王莹', 'address': '江西省博市房山太原路N座 615506', 'birthdate': datetime.date(1968, 2, 7)};

示例代码:https://github.com/JustDoPython/python-examples/tree/master/doudou/2021-01-10-fake-data

解析MySQL Binlog文件

binlog2sql 是通过分析 MySql 数据库的 binlog 文件,从中解析出需要执行的 sql 语句的。

那么使用时需要提供一些必要的参数,其中重要的有数据库服务器链接信息,需要分析的 binlog 文件名等,

还可以指定解析的起始和结束位置,以及开始和结束时间。

安装

binlog2sql 使用 Python 开发,所以需要 Python 环境,可参考 Python 环境搭建

将 binlog2sql 用 git 克隆的本地,GitHub 上的地址是: https://github.com/danfengcao/binlog2sql.git

git clone https://github.com/danfengcao/binlog2sql.git 

通过 binlog2sql 目标下的 requirements.txt 安装依赖包

提示:推荐在 Python 虚拟环境中安装,创建虚拟环境可参考 Python 虚拟环境 看这一篇就够了

pip install -r requirements.txt 

一切顺利的话,很快就可完成安装。

命令行进入 binlog2sql 代码目录下测试一下

> python binlog2sql.py

usage: binlog2sql.py [-h HOST] [-u USER] [-p [PASSWORD ...]] [-P PORT] [--start-file START_FILE] [--start-position START_POS] [--stop-file END_FILE] [--stop-position END_POS]
                     [--start-datetime START_TIME] [--stop-datetime STOP_TIME] [--save-as SAVE_AS] [--stop-never] [--help] [-d [DATABASES ...]] [-t [TABLES ...]] [--only-dml]
                     [--sql-type [SQL_TYPE ...]] [-K] [-B] [--back-interval BACK_INTERVAL]

Parse MySQL binlog to SQL you want

...<省略>...

由于没加任何参数,所以打印出使用说明,那说明安装正常了。

恢复被删数据

假如库表 tb_user 中的数据如下:

+----+--------+---------------------+
| id | name   | createtime          |
+----+--------+---------------------+
|  1 | 张三   | 2021-01-10 00:04:33 |
|  2 | 李四   | 2021-01-10 00:04:48 |
|  3 | 王五   | 2021-04-23 20:25:00 |
|  4 | 赵六   | 2021-06-04 11:21:23 |
+----+--------+---------------------+

这时不小心执行了一个删操作,将数据误删了

delete from tb_user 

如何恢复呢?

我们看一下数据库的日志情况

show master status; 

会看到类似这样的结果

+------------------+-----------+
| File             | Position  |
+------------------+-----------+
| mysql-bin.000002 |     13136 |
+------------------+-----------+

注意:只有 MySql 数据库打开了日志记录功能,才能查询到,打开日志功能请参考 binlog日志开启和使用

可以看出,目前日志记录在文件 mysql-bin.000002 中,当前最新的记录位置是 12546 行

假如当时误操作的时间是上午 11点半左右(可能着急吃饭,没注意),那么预估一个时间范围,比如 11点25 到 11点35,看看一下当时的操作:

python binlog2sql -h127.0.0.1 -P3306 -uadmin -p'admin' -dtest -t tb_user --start-file='mysql-bin.000002' --start-datetime='2021-06-04 11:25:00' --stop-datetime='2021-06-04 11:35:00'

输出为:

INSERT INTO `test`.`tb_user`(`createtime`, `id`, `name`) VALUES ('2021-06-04 11:21:23', 4, '李四'); #start 12317 end 12487 time 2021-06-04 11:21:23 DELETE FROM `test`.`tb_user` WHERE `createtime`='2021-01-10 00:04:33' AND `id`=1 AND `name`='张三' LIMIT 1; #start 12728 end 12829 time 2021-06-04 11:27:32 DELETE FROM `test`.`tb_user` WHERE `createtime`='2021-01-10 00:04:48' AND `id`=2 AND `name`='李四' LIMIT 1; #start 12728 end 12829 time 2021-06-04 11:27:32 DELETE FROM `test`.`tb_user` WHERE `createtime`='2021-04-23 20:25:00' AND `id`=3 AND `name`='王五' LIMIT 1; #start 12728 end 12829 time 2021-06-04 11:27:32 DELETE FROM `test`.`tb_user` WHERE `createtime`='2021-06-04 11:21:23' AND `id`=4 AND `name`='赵六' LIMIT 1; #start 12728 end 12829 time 2021-06-04 11:27:32 

可以看出,第二行开始到第五行为删除语句,查看语句最后的起始和结束位置 start 12728 end 12829

即 binlog 中,删除执行的位置在 12728-12829 之间,于是锁定精确位置,生成回滚语句:

python binlog2sql -h127.0.0.1 -P3306 -uadmin -p'admin' -dtest -t tb_user --start-file='mysql-bin.000002' --start-position=12728 --stop-position=12829 -B

注意参数 -B,意思是生成回滚 SQL,即生成的是撤销之前操作的语句

输出为:

INSERT INTO `test`.`tb_user`(`createtime`, `id`, `name`) VALUES ('2021-06-04 11:21:23', 4, '赵六'); #start 12728 end 12829 time 2016-12-13 20:28:05
INSERT INTO `test`.`tb_user`(`createtime`, `id`, `name`) VALUES ('2021-04-23 20:25:00', 3, '王五'); #start 12728 end 12829 time 2016-12-13 20:28:05
INSERT INTO `test`.`tb_user`(`createtime`, `id`, `name`) VALUES ('2021-01-10 00:04:48', 2, '李四'); #start 12728 end 12829 time 2016-12-13 20:28:05
INSERT INTO `test`.`tb_user`(`createtime`, `id`, `name`) VALUES ('2021-01-10 00:04:33', 1, '张三'); #start 12728 end 12829 time 2016-12-13 20:28:05

从输出的语句来看,顺序是删除的倒序,而且已经将原来的 delete 语句改为了 insert 语句,也就是原来操作的逆操作

如果确认语句没问题,执行生成的语句就可以了

是不是既方便又高效呢?

解析 SQL

binlog2sql 功能强大,使用起来也很方便,看看其他功能吧。

作为一个命令行工具,功能都体现在参数里,可分为 解析模式、解析目标、解析范围三部分。

解析模式

binlog2sql 支持两个解析模式,默认的是单次解析,即运行一次解析一次,

还可以支持持续解析,即不间断地从目标数据库地 binlog 中解析出 sql 来,持续解析通过参数 --never-stop 开启,

开启之后,线程不会退出,一直处于运行状态,会自动判断 binlog 的变化,对变化部分增量式解析。

这种模式可以用于数据库同步,不过生产上使用前,最好考虑各种异常情况,比如重启,网络中断等情况。

参数 -K--no-primany-key 表示的去除 INSERT 语句中的主键,这个在数据汇总的场景下很方便,可以避免多个数据源中主键冲突的问题。

参数 -B--flashback,表示回滚模式,在上面的例子中展示过,即会解析成逆操作的 sql 语句。

在回滚模式下,每生成一千条 SQL 语句会加一个 SLEEP 语句,是为以免数据执行时产生拥堵,默认为 1 秒,可以通过 --back-interval 参数来设置,

例如 --back-interval 2 表示暂停 2 秒。

解析目标

MySql 设置 binlog 时可以指定记录哪个库,以及哪些表,即目标。

那么用 binlog2sql 也可以指定解析目标。

参数 -d--databases 用于指定数据库,如果多个库,用空格分隔,例如 -d db1 db2

参数 -t--tables 用于指定库表,多个库表用空格分隔,例如 -t tb1 tb2

如果指定解析目标不仅效率更高,而且分析和执行解析的结果也更方便。

解析范围

范围包括 binlog 文件范围时间范围 以及 行范围,例如前面例子中用到了 时间范围行范围

文件范围--start-file--stop-file 参数来指定,只需要提供 binlog 文件名即可,不需要写全路径,这是因为,binlog2sql 会自动根据目标服务器配置读取 binlog 文件;

时间范围--start-datetime--stop-datetime 参数来指定,时间格式为 %Y-%m-%d %H:%M:%S

行范围--start-position--stop-position 参数来指定,也可以简写为 --start-pos--end-pos

深入了解

binlog2sql 不仅是一个实用的工具,而且也是个研究和学习的好例子。

只有不到 500 行代码,很容易阅读;

阅读源码,不仅能深入了解其实现原理,而且还可以学习到很多好用法。

实现原理

binlog2sql 的原理是,利用 pymysql 从目标服务器上获取 binlog 信息,然后锁定范围,使用 pymysqlreplication 解析 binlog 文件,最后,得到需要解析出的 sql 语句。

在这基础上,做了一些功能性扩展,比如解析范围,解析模式等,相对来说比较简单,很容易看懂。

命令行参数

编程时处理命令行参数是机械而繁琐的,特别是有不同性质的性质和别名的参数

binlog2sql 中利用了 argparse 模块,

argparse 模块可以让人轻松编写用户友好的命令行接口。程序定义它需要的参数,argparse 可以从 sys.argv 解析出提供的命令行参数。而且 argparse 模块还会自动生成帮助和使用手册,并在用户给程序传入无效参数时报出错误信息。

很容易就能编程高大上的命令行程序接口,再也不用为很 low 的程序接口发愁了。

文件处理上下文

binlog2sql 在回滚模式(即提供了参数 -B)中,使用了一个临时文件记录解析出来的 SQL 语句,并且在完成之后删除。

一般来说,在完成后主动删除文件即可,不过如果能利用 with 块的资源回收功能就更好了。

查看源码,会看到一个创建文件写法:

@contextmanager
def temp_open(filename, mode):
    f = open(filename, mode)
    try:
        yield f
    finally:
        f.close()
        os.remove(filename)

@contextmanager 指示器可以将一个生成器,作为一个上下文管理器,

那么:

with 声明部分,会执行前会执行 yield 语句之前的部分

with 范围内,会执行 yield 语句,即返回一个需要后续处理的对象,比如文件,后续处理是关闭

with 执行完成前,会执行 yield 语句之后的代码

那么这段代码的意义就是,当文件使用完成后,关闭文件,并且删除掉。

使用方式为:

with temp_open(tmp_file, "w") as f_tmp
    ...
    f_tmp.write(sql + '\n')
    ...

这样无论如何只要 with 块执行完,文件就会被删除,不用担心忘记,是不是很优雅?

原文 :http://www.justdopython.com/2021/06/04/binlog2sql/

计算代码执行时间

def timer(func):
    def wrapper(*args, **kwargs):
        start = time.time()
        res = func(*args, **kwargs)
        print('共耗时约 {:.2f} 秒'.format(time.time() - start))
        return res

    return wrapper


@timer
def gen():
    """ 使用生成器批量写入数据 """
    helpers.bulk(es, action)

容器镜像构建

python 源码

目录结构

image-20210724110225875

/task/static/socket-2.0.io.min.js 官网下载(已被墙):https://socket.io/blog/socket-io-2-0-0/

/task/jquery-3.6.0.js 官网下载:https://jquery.com/

/task/templates/hello.html




    Flask-SocketIO Test
    
    
    
    
    


WebSokect

/task/app.py

from flask import Flask, render_template, session, request
from flask_socketio import SocketIO, emit
import random, time

app = Flask(__name__)
app.config['SECRET_KEY'] = 'secret!'
socketio = SocketIO(app)

@app.route('/', methods=['POST', 'GET'])
def index():
    return render_template('hello.html')


@socketio.on('client_event')
def client_msg(msg):
    while True:
        time.sleep(3)
        random_ = random.randint(0, 100)
        print(random_)
        emit('server_response', {'data': random_})


@socketio.on('connect_event')
def connected_msg(msg):
    emit('server_response', {'data': msg['data']})


if __name__ == '__main__':
    socketio.run(app, host='0.0.0.0',debug=True)

/task/requirements.txt

Flask~=2.0.1
Flask-SocketIO==4.3.1
python-engineio==3.13.2
python-socketio==4.6.0

编写镜像文件

FROM python:3.8.6
WORKDIR /opt/task
COPY ./task /opt/task
RUN pip3 install -i https://pypi.douban.com/simple -r requirements.txt
EXPOSE 5000
CMD ["flask", "run","host","0.0.0.0"]

临时用,-i 指定下载源,国内源:

  • 清华:https://pypi.tuna.tsinghua.edu.cn/simple
  • 阿里云:http://mirrors.aliyun.com/pypi/simple/
  • 中国科技大学 https://pypi.mirrors.ustc.edu.cn/simple/
  • 华中理工大学:http://pypi.hustunique.com/
  • 山东理工大学:http://pypi.sdutlinux.org/
  • 豆瓣:http://pypi.douban.com/simple/

永久修改(可选)

Linux下,修改 ~/.pip/pip.conf (没有就创建一个文件夹及文件。文件夹要加“.”,表示是隐藏文件夹)

内容如下:

[global]
index-url = https://pypi.tuna.tsinghua.edu.cn/simple
[install]
trusted-host=mirrors.aliyun.com

windows下,直接在user目录中创建一个pip目录,再新建文件pip.ini。
(例如:C:\Users\你的用户名\AppData\Roaming\pip)内容同上。

构建

docker image build -t flask-websocket:1.0 .

运行

docker run -d --name flask-websocket -p 5000:5000 flask-websocket:1.0

Python调用ES

第一种方式

import time
import ArticleGenerated
from mimesis import Person, Address
from elasticsearch import Elasticsearch


class ElasticObj:
    def __init__(self, index_name, index_type, ip, port):
        self.index_name = index_name
        self.index_type = index_type
        self.ip = ip
        self.port = port
        self.es = Elasticsearch([ip], port=port)
        #使用账号密码
        #es = Elasticsearch(['localhost:9200'], http_auth="username:password")

    def GET(self):
        res = self.es.search(index=self.index_name, filter_path=["hits.hits._*"])
        print(res)

    """数据生成"""
    def POST(self):
        person = Person('zh')
        address = Address('zh')
        title = address.continent()
        user = person.surname() + person.name()
        desc = ArticleGenerated.Generated(title)
        now_time = time.strftime('%Y-%m-%d %H:%M')

        index_data = {"title": title, "user": user, "desc": desc, "time": now_time}
        res = self.es.index(index=self.index_name, document=index_data)
        print(res['result'])


    def PUT(self):
        pass

    def DELETE(self):
        pass


if __name__ == '__main__':
    obj = ElasticObj('test2', '_doc', '192.168.111.140', '9201')
 
    obj.POST()
    obj.es.indices.refresh(index=obj.index_name)

第二种方式

bulk API可以在单个请求中一次执行多个操作(index,udpate,create,delete),使用这种方式可以极大的提升索引性能。

在这里我们使用elasticsearch模块的helpers,helpers是bulk的帮助程序,是对bulk的封装。有三种方式bulk(),streaming_bulk(),parallel_bulk(),都需要接受Elasticsearch类的实例和一个可迭代的actions.actions里面的操作可以有多种。

官网:https://elasticsearch-py.readthedocs.io/en/master/helpers.html

import random
from elasticsearch import Elasticsearch
from elasticsearch import helpers

es=Elasticsearch(hosts='http://localhost',port=9200)
levels=['info','debug','warn','error']
actions=[]
for i in range(100):
    level=levels[random.randrange(0,len(levels))]
    action={'_op_type':'index',#操作 index update create delete  
            '_index':'log_level',#index
            '_type':'doc',  #type
            '_source':{'level':level}}
    actions.append(action)
#使用bulk方式
helpers.bulk(client=es,actions=actions)
#streaming_bulk与parallel_bulk类似  需要遍历才会运行
#都可以设置每个批次的大小,parallel_bulk还可以设置线程数  
for ok,response in helpers.streaming_bulk(es,actions):
        if not ok:
            print(response)

参考:

一次性插入十万条数据

import time
from elasticsearch import Elasticsearch
from elasticsearch import helpers

es = Elasticsearch()

def timer(func):
    def wrapper(*args, **kwargs):
        start = time.time()
        res = func(*args, **kwargs)
        print('共耗时约 {:.2f} 秒'.format(time.time() - start))
        return res

    return wrapper
@timer
def gen():
    """ 使用生成器批量写入数据 """
    action = ({
        "_index": "s2",
        "_type": "doc",
        "_source": {
            "title": i
        }
    } for i in range(100000))
    helpers.bulk(es, action)

if __name__ == '__main__':
    # create_data()
    # batch_data()
    gen()

参考:

查询数据

get获取

res = es.get(index="my-index", doc_type="test-type", id=01)
es.get(index='indexName', doc_type='typeName', id='idValue')

搜索所有数据

es.search(index="my_index",doc_type="test_type")
# 或者
body = {
    "query":{
        "match_all":{}
    }
}
es.search(index="my_index",doc_type="test_type",body=body)

term与terms

body = {
    "query":{
        "term":{
            "name":"python"
        }
    }
}
# 查询name="python"的所有数据
es.search(index="my_index",doc_type="test_type",body=body)
terms

body = {
    "query":{
        "terms":{
            "name":[
                "python","android"
            ]
        }
    }
}
# 搜索出name="python"或name="android"的所有数据
es.search(index="my_index",doc_type="test_type",body=body)

match与multi_match

# match:匹配name包含python关键字的数据
body = {
    "query":{
        "match":{
            "name":"python"
        }
    }
}
# 查询name包含python关键字的数据
es.search(index="my_index",doc_type="test_type",body=body)

multi_match:在name和addr里匹配包含深圳关键字的数据

body = {
    "query":{
        "multi_match":{
            "query":"深圳",
            "fields":["name","addr"]
        }
    }
}
# 查询name和addr包含"深圳"关键字的数据
es.search(index="my_index",doc_type="test_type",body=body)
#切片式查询

body = {
    "query":{
        "match_all":{}
    }
    "from":2    # 从第二条数据开始
    "size":4    # 获取4条数据
}
# 从第2条数据开始,获取4条数据
es.search(index="my_index",doc_type="test_type",body=body)

范围查询

body = {
    "query":{
        "range":{
            "age":{
                "gte":18,       # >=18
                "lte":30        # <=30
            }
        }
    }
}
# 查询18<=age<=30的所有数据
es.search(index="my_index",doc_type="test_type",body=body)

只需要获取_id数据,多个条件用逗号隔开

es.search(index="my_index",doc_type="test_type",filter_path=["hits.hits._id"])

获取所有数据

es.search(index="my_index",doc_type="test_type",filter_path=["hits.hits._*"])

获取数据量

es.count(index="my_index",doc_type="test_type")

删除数据

delete:删除指定index、type、id的文档

es.delete(index='indexName', doc_type='typeName', id='idValue')

条件删除

delete_by_query:删除满足条件的所有数据,查询条件必须符合DLS格式

query = {'query': {'match': {'sex': 'famale'}}}# 删除性别为女性的所有文档

query = {'query': {'range': {'age': {'lt': 11}}}}# 删除年龄小于11的所有文档

es.delete_by_query(index='indexName', body=query, doc_type='typeName')

更新数据

条件更新

update_by_query:更新满足条件的所有数据,写法同上删除和查询

批量写入、删除、更新

delete_by_query:删除满足条件的所有数据,查询条件必须符合DLS格式

query = {'query': {'match': {'sex': 'famale'}}}# 删除性别为女性的所有文档

query = {'query': {'range': {'age': {'lt': 11}}}}# 删除年龄小于11的所有文档

es.delete_by_query(index='indexName', body=query, doc_type='typeName')

参考

https://elasticsearch-py.readthedocs.io/en/v7.15.1/api.html#global-options API Documentation

Redis操作

几种连接方式介绍

python要调用redis的时候,需要先安装redis模块,有两个方法。第一个方法就是pip install redis,第二个方法就是easy_install redis,模块装完之后,就可以创建redis连接了。

在这个库中有两个类Redis和StrictRedis来实现Redis的命令操作。Redis是StrictRedis的子类,主要功能是向后兼容旧版本库里的几个方法。在这里使用官方推荐的StrictRedis。

from redis import StrictRedis

# 连接redis需要IP地址,运行端口,数据库,密码。如果本地redis没有密码可以直接使用默认值
conn = StrictRedis(host="localhost", port=6379, db=0, password=None
print(type(conn))
# 连接成功之后调用实例对象的set方法设置值
conn.set("name", "Yang")
# 调用get方法获取值
print(conn.get("name"))

关系型数据库都有一个连接池的概念:对于大量redis连接来说,如果使用直接连接redis的方式的话,将会造成大量的TCP的重复连接,所以,就引入连接池来解决这个问题。在使用连接池连接上redis之后,可以从该连接池里面生成连接,调用完成之后,该链接将会返还给连接池,供其他连接请求调用,这样将减少大量redis连接的执行时间,那么使用StrictRedis的连接池的实现方式如下:

pool = redis.ConnectionPool(host=host, port=6379, password=password)
r = redis.StrictRedis(connection_pool=pool

或者使用pipeline(管道),通过缓冲多条命令,然后一次性执行的方法减少服务器-客户端之间TCP数据库包,从而提高效率,方法如下:

#接上文
pipe = r.pipeline()
#插入数据
pipe.hset("hash_key","leizhu900516",8)
#Pipeline>>

pipe.hset("hash_key","chenhuachao",9)
#Pipeline>>

pipe.hset("hash_key","wanger",10)
#Pipeline>>
pipe.execute()
#[1L, 1L, 1L]

批量读取数据的方法如下:

pipe.hget("hash_key","leizhu900516")
#Pipeline>>
pipe.hget("hash_key","chenhuachao")
#Pipeline>>
pipe.hget("hash_key","wanger")
#Pipeline>>
result = pipe.execute()
print result
#['8', '9', '10']   #有序的列表

pipeline的命令可以写在一起,如p.set('hello','redis').sadd('faz','baz').incr('num').execute(),其实它的意思等同于是:

>>> p.set('hello','redis')
>>> p.sadd('faz','baz')
>>> p.incr('num')
>>> p.execute()
[True, 1, 1]

利用pipeline取值3500条数据,大约需要900ms,如果配合线程or协程来使用,每秒返回1W数据是没有问题的,基本能满足大部分业务。

插入大量数据

if __name__ == '__main__':
    pool = redis.ConnectionPool(host='120.27.243.25', port=6379, decode_responses=True, password='^du5PXq}hn>r]r6T}LNw')
    r = redis.Redis(connection_pool=pool)

    # print(r.get('food'))  # 取出键food对应的值
    list_key = []
    list_value = []
    for x in range(1000000):
        key = ''.join(random.sample(string.ascii_letters + string.digits, 8))
        value = ''.join(random.sample(string.ascii_letters + string.digits, 8))
        list_key.append(key)
        list_value.append(value)
    p = r.pipeline()  # 创建一个管道,现在本地保存命令,插入执行完后,一次性插入到redis
    for k, v in zip(list_key, list_value):
        #print(k, v)
        # 下面的 SET 命令被缓冲
        p.set(k, v)
    # EXECUTE 调用将所有缓冲的命令发送到服务器,返回
    p.execute()  # 执行管道内命令

官文:https://pypi.org/project/redis/

插入JSON

import json
from mimesis import Person
from redis import StrictRedis

if __name__ == '__main__':
    person = Person('zh')
    person_json = {
        "name": person.surname() + person.name(),
        "Sex": person.sex(),
        "academic degree": person.academic_degree()
    }
    conn = StrictRedis(host="localhost", port=6379, db=0, password=None)
    conn.lpush("testjson", json.dumps(person_json))

操作

更多: https://cloud.tencent.com/developer/article/1151834

键操作

exists(name)        判断一个键是否存在               redis.exists("name")
delete(name)        删除一个键                       redis.delete("name")
type(name)          判断键的类型                     redis.type("name")
keys(pattern)       获取所有符合规则的键             redis.keys("n*")
randomkey()         获取随机的一个键                 redis.randomkey(self)   
rename(src, dst)    重命名键                         redis.rename("name", "names")
dbsize()            获取当前数据库中键的数目          redis.dbsize(self)
expire(name, time)  设定键的过期时间,单位:秒        redis.expire("name", 10)
ttl(name)           获取键的过期时间,单位:秒        redis.ttl("name")      -1表示永久
move(name, db)      将键移动到其他数据库              redis.move("name", 1)
flushdb()           删除当前选择数据库中的所有键      redis.flushdb(self)
flushall()          删除所有数据库中的所有键          redis.flushall(self)

字符串操作

set(name, value)        将键name的值设置为`value`      redis.set('name', 'Yang')
get(name)               获取键名为name的值             redis.get('name')
getset(name, value)     将键name的值设置为value并以原子方式返回键name的旧值。   redis.getset("name", "Yu")
mget(keys, *args)       返回多个键对应的value           redis.mget(['name', 'names'])
setnx(name, value)      不存在则创建,存在则不生效       redis.setnx("newname", "Yang")
setex(name, time, value)    设置一个有效期为1秒的值      redis.setex("name", 1, "Yang")
setrange(name, offset, value)   为指定键的值拼接字符串   redis.setrange("name", num, "value")  offset为偏移量也就是指针  
mset(mapping)           批量赋值                        redis.mset({"name": "value"})
msetnx(mapping)         键都不存在时批量赋值             redis.msetnx({"name": "value"})
incr(name, amount=1)    键不存在则创建并设置值为1,存在则值增长1     redis.incr("score": 1)
decr(name, amount=1)    与上一个方法相反、减值操作
append(key, value)      为key的值追加value               redis.append("name", "Yu")
substr(name, start, end=-1)   截取键name值的字符串      redis.substr("name", 1, 3)      end=-1表示截取到末尾
getrange(key, start, end)     与上一个方法相同          redis.getrange("name", 1, 3)

列表操作

rpush(name, *value)     在name键的列表值末尾追加多个元素      redis.rpush("name", 1,2,3)
lpush(name, *value)     在name键的列表值头部推送多个元素      redis.lpush("name", 0)
llen(name)              返回name键的列表值长度                redis.llen("name")
lrange(name, start, end)    返回name键的切片列表值            redis.lrange("name", 1, 3)
ltrim(name, start, end)     与上一个方法相同
lindex(name, index)     返回name键的列表中index位置的元素     redis.lindex("name", 2)
lset(name, index, value)    给name键的列表中index位置的元素赋值,越界则报错    redis.lset("name", 2, "Yu")
lrem(name, count, value)    删除name键的列表中count个值为value的元素     redis.lrem("name", 2, "Yang")
lpop(name)              返回并删除name键的列表中首元素        redis.lpop("name")
rpop(name)              与上一个方法相反
blpop(keys, timeout=0)  返回并删除name键的列表中首元素       如果列表为空则阻塞,timeout=0表示一直阻塞
brpop(keys, timeout=0)  与上一个方法相反                      redis.brpop("keys")
rpoplpush(src, dst)     将src键的列表中首元素取出添加到dst键列表中的头部     redis.rpoplpush("name", "name1")
lrange("mylist", 0, -1) #取出全部键值对

集合操作

sadd(name, *values)     向键为name的集合中添加多个元素     redis.sadd("tags", "Yang", "Yu")
srem(name, *values)     从键为name的集合中删除多个元素     redis,srem("tags", "Yu")
spop(name)              在集合name中随机删除并返回删除的元素     redis.spop("tags")
smove(src, dst, value)  从src集合中移除元素添加到dst集合中       redis.smove("tags", "tags2", "Yang")
scard(name)             返回集合name中的元素个数           redis.scard("tags")
sismember(name, value)  判断value是否在集合name中          redis.sismember("tags", "Yang")
sinter(keys, *args)     返回给定键keys中多个集合的交集     redis.sinter(["tags", "tags2"])
sinterstore(dest, keys, *args)   求交集并保存到dest集合中  redis.sinterstore("intag", ["tags", "tags2"])
sunion(keys, *args)     返回所有给定键keys的并集           redis.sunion(["tags", "tags2"])
sunionstore(dest, keys, args)    求并集并保存到dest集合中  redis.sunionstore("intags", ["tags", "tags2"])
sdiff(keys, *args)      返回所有给定键集合的差集           redis.sdiff(["tags", "tags2"])
sdiffstore(dest, keys, *args)    求差集并保存到dest集合中  redis.sdiffstore("intags", ["tags", "tags2"])
smembers(name)          返回键为name的集合中所有元素       redis.smembers("tags")
srandmember(name)       随机返回键为name集合中的一个元素并不删除  redis.srandmember("tags")

有序集合操作

zadd(name, *args, **kwargs) 向键为name的zset中添加多个元素     redis.zadd("tags", 91, "Jim", 80, "Bob")
zrem(name, *values)         删除键为name的zset中的元素         redis.zrem("tags", "Bob")
如果在键为name的zset中已经存在元素value,则将该元素的score增加amount,
否则向该集合中添加该元素,其score的值为amount
zincrby(name, value, amount=1)                                 redis.zincrby("tags", "Jim". -9)
zrank(name, value)          返回键为name中score从小到大的zset排名   redis.zrank("tags", "Jim")           
zrevrank(name, value)       与上一个方法相反
返回键为name的zset中index从start到end的所有元素
zrevrange(name, start, end, withscores=False)       redis.zrevrange("tags", 0, 2)  
返回键为name的zset中score在给定区间的元素, num:个数,withscores:是否带score
zrangebyscore(name, min, max, start=None, num=None, withscores=False)   redis.zrangebyscore("tags", 80, 90)
zcount(name, min, max)      返回键为name的zset中score在给定区间的元素个数   redis.zcount("tags", 80, 90)
zcard(name)                 返回键为name的zset中元素个数        redis.zcard("tags")
删除键为name的zset中排名在给定区间的元素
zremrangebyrank(name, min, max)     redis.zremrangebyrank("tags", 0, 59)
删除键为name的zset中score在给定区间的元素
zremrangebyscore(name, min, max)    redis.zremrangebyscore("tags", 80, 90)

散列操作

hset(name, key, value)      向键为name的散列表中添加映射            redis.hset("price", "apple", 2)
hsetnx(name, key, value)    不存在则向键为name的散列表中添加映射    redis.hsetnx("price", "orange", 1)
hget(name, key)             返回键为name的散列表中key对应的值       redis.hget("price", "apple")
hmget(name, keys, *args)    返回键为name的散列表中各个key对应的值   redis.hmget("price", ["apple", "orange"])
hmset(name, mapping)        向键为name的散列表中批量添加映射        redis.hmset("price", {"banana": 2, "tomato": 1})
hincrby(name, key, amount=1) 向键为name的散列表中对应key的值增加    redis.hincrby("price", "apple", 1)
hexists(name, key)          判断键为name的散列表中是否存在键为key的映射   redis.hexists("price", "tomato")
hdel(name, *keys)           删除键为name的散列中对应key             redis.hdel("price", "banana")
hlen(name)                  从键为name的散列表中获取映射个数        redis.hlen("price")
hkeys(name)                 从键为name的散列表中获取所有映射键名    redis.hkeys("price")
hvals(name)                 从键为name的散列表中获取所有映射键值    redis.hvals("price")
hgetall(name)               从键为name的散列表中获取所有映射键值对  redis.hgetall("price")

Python配置文件

image-20211123150436520

import configparser

cf = configparser.ConfigParser()
cf.read("E:\Crawler\config.ini")  # 读取配置文件,如果写文件的绝对路径,就可以不用os模块

#获取指定Email section的host键
print(cf.get('Email', 'host'))

secs = cf.sections()  # 获取文件中所有的section(一个配置文件中可以有多个配置,如数据库相关的配置,邮箱相关的配置,
                      # 每个section由[]包裹,即[section]),并以列表的形式返回
print(secs)

options = cf.options("Mysql-Database")  # 获取某个section名为Mysql-Database所对应的键
print(options)

items = cf.items("Mysql-Database")  # 获取section名为Mysql-Database所对应的全部键值对
print(items)

host = cf.get("Mysql-Database", "host")  # 获取[Mysql-Database]中host对应的值
print(host)

参考

https://blog.csdn.net/liaodaoluyun/article/details/82346086 Elasticsearch的bulk用法(python)

http://www.justdopython.com/archives/

https://rorschachchan.github.io/2018/03/26/使用python调用redis的基本操作/ 使用python调用redis的基本操作