开源数据流框架之luigi介绍

一,介绍

Luigi 是一个 Python 模块,可以帮你构建复杂的批量作业管道。处理依赖决议、工作流管理、可视化展示等等,内建 Hadoop 支持。它也被Foursquare,Stripe,华尔街日报,Groupon和其他知名企业使用。

Luigi是基于代码的,而不是基于GUI或声明式的,包含Python中的所有内容(包括依赖关系图)。用户界面(UI)允许您搜索,过滤或监视每个任务的状态。您还可以查看该工作流程,以查看依赖关系图上的哪些任务已完成,哪些尚未运行。


image.png

Luigi轻便简洁, 能够快速实现并行化的工作流。


image.png

二,安装和基本概念

1,安装

pip install luigi

2, 基本概念

luigi 有三个重要概念,Task, Target 和parameter

Task

每一个任务都是一个Task,以class的形式存在,继承luigi.Task。需要重载requires()、run()、output()方法。
其中requires()是任务入口程序,指定任务依赖的上游输入;run()是任务在该节点具体实现的流程;output()是任务的出口,把该节点执行完之后的结果输出到下游。


image.png

Target

广义地讲,Target可对应为磁盘上的文件,或HDFS上文件,或checkpoint点,或数据库等。对于Target来说,唯一需要实现的方法为exists,返回为True表示存在,否则不存在返回为False. 在实际应用时,写一个Target子类是很少需要用到的。直接使用开箱即可用的LocalTarget及 hdfs.HdfsTarget类就够用了。Luigi提供了Gzip支持,通过参数format=format.Gzip即可。

自定义Target的示例

# 实现判断数据库中test_table表是否被创建的Target
import luigi
import MySQLdb

class ConnectDB:

    def __init__(self, host="128.0.0.1", user="root", passwd="root", db="test", port=3306):
        self.conn = MySQLdb.connect(host, user, passwd, db, port=port, charset='utf8')

    def __enter__(self):
        return self.conn

    def __exit__(self, exc_type, exc_val, exc_tb):
        self.conn.close()

    def connect(self):
        return self.conn

class MysqlTarget(luigi.Target):
    def __init__(self, primer_table):
        self.primer_table = primer_table

    def exists(self):
        with ConnectDB() as conn:
            sql = "SELECT TABLE_NAME FROM INFORMATION_SCHEMA.TABLES WHERE TABLE_NAME = '{}'".format(self.test_table)
            cur = conn.cursor()
            cur.execute(sql)
            table = cur.fetchone()
            cur.close()
            return True if table else False

parameter

parameter等效于luigi为task类创建构造函数,Luigi中提供了不同类型的parameter,例如DateParameter,DateIntervalParameter,IntParameter,FloatParameter等等。
python不是一个静态类型的语言,你不需要指定参数的类型,你可以直接使用基类Parameter。

class TaskA(luigi.Task):
    x = luigi.Parameter()
    y = luigi.IntParameter(default=0)

三,示例代码

在这个例子中我们定义PrintNum任务创建一串数字写入到文件中,其他平方SquareNum,求倍数AddNum和立方CubeNum任务依赖这个文件,把自己的计算结果写如到自己对应的文件中。

#!/usr/bin/python
# _*_ coding:utf-8 _*_

import luigi
import time
import uuid

UID = str(uuid.uuid4())

# 生成1, n个数字,写入到number_up_{}.txt中
class PrintNum(luigi.Task):
    n = luigi.IntParameter(default=5)
    lt = list()

    def requires(self):
        return []

    def output(self):
        return luigi.LocalTarget("number_up_{}.txt".format(UID))

    def run(self):
        with self.output().open("w") as f:
            for i in range(1, self.n):
                self.set_status_message("Progress: %d" % i)
                self.set_progress_percentage(i)
                f.write("{}\n".format(i))

# 依赖PrintNum,对数字求平方后写入到squares_{}.txt中
class SquareNum(luigi.Task):
    n = luigi.IntParameter(default=10)

    def requires(self):
        return [PrintNum()]

    def output(self):
        return luigi.LocalTarget("squares_{}.txt".format(UID))

    def run(self):
        with self.input()[0].open() as fin, self.output().open("w") as fout:
            for line in fin:
                n = int(line.strip())
                list(range(10**7))
                # time.sleep(1)
                out = n * n
                fout.write("{}:{}\n".format(n, out))

# 依赖PrintNum,对数字求倍数后写入到add_{}.txt中
class AddNum(luigi.Task):
    n = luigi.IntParameter(default=10)

    def requires(self):
        return [PrintNum(n=self.n)]

    def output(self):
        return luigi.LocalTarget("add_{}.txt".format(UID))

    def run(self):
        with self.input()[0].open() as fin, self.output().open("w") as fout:
            for line in fin:
                n = int(line.strip())
                list(range(10**7))
                # time.sleep(1)
                out = n + n
                fout.write("{}:{}\n".format(n, out))

# 依赖PrintNum,对数字求立方后写入到cube_{}.txt中
class CubeNum(luigi.Task):
    n = luigi.IntParameter(default=10)

    def requires(self):
        return[PrintNum(n=self.n)]

    def output(self):
        return luigi.LocalTarget("cube_{}.txt".format(UID))

    def run(self):
        with self.input()[0].open() as fin, self.output().open("w") as fout:
            for line in fin:
                n = int(line.strip())
                list(range(10**7))
                # time.sleep(1)
                out = n ** 3
                fout.write("{}:{}\n".format(n, out))

if __name__ == "__main__":
    #构建任务
    luigi.build([SquareNum(n=40), AddNum(n=40), CubeNum(n=40)], local_scheduler=True)
最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 226,097评论 6 523
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 97,198评论 3 410
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 173,602评论 0 370
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 61,750评论 1 304
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 70,684评论 6 404
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 54,151评论 1 317
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 42,349评论 3 433
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 41,430评论 0 282
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 47,986评论 1 328
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 39,969评论 3 351
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 42,071评论 1 359
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 37,667评论 5 352
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 43,358评论 3 342
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 33,757评论 0 25
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 34,944评论 1 278
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 50,684评论 3 384
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 47,123评论 2 368