Cloud Data Fusionの機能でつまづいた部分について
Google Cloud Platform Advent Calendar 2021の22日目カレンダー2の記事です。
最近Cloud Data Fusionをはじめて触りました。GUIベースで結構感覚的に処理ができるので、結構いいツールだと思います。
そして今回はCloud Data Fusionで説明が必要だなあと感じた機能について説明します。
Cloud Data Fusionについての説明や、インスタンスの作成などは公式ドキュメントを参考にしてください。
今回使用するサンプルのテーブルとデータは下記を使用します。

RawDenormalizer
BigQueryのテーブルの行列をPIVOT(行と列の入れ替え)したいときはRowDenormalizerを使います。
ただ、dateとsub_keyをPIVOTのkeyとして、nameとvalをPIVOT したいとき(下記のようなテーブルを想定)が面倒くさいです。理由はRowDenormalizerがkeyを1つしか指定できないためです。
そのため、Wranglerを使ってKeyにしたい列をmergeしてからRawDenormalizerをする必要があります。
Wranglerでの列Joinはこのように行います

そして、RowDenormalizerでこのように指定します。

RowDenormalizerをした結果はこのようなデータになります。

Join
SQLで使うJoinは下記図の上のように使います。
ただ、縦結合、UNIONでデータをJoinさせたい場合は、下記図の下のように使います。

テーブル出力
パイプラインを実行するたびにテーブルの中身を上書きしたい場合は、設定を変更する場合があります。Truncate TableをTrueにします。

ちなみに、デフォルトではパイプラインを実行するたびにテーブルにInsertする処理が設定されています(Insert以外にUpdate, Upsertを選択できる)。
requests.Session をシングルトンで書いて使い回す
シングルトンの実装をしたrequests.Sessionクラスを使うことで、毎回同じ識別子のクラスを使い回すことが出来ます。
そうすることで、リクエストを飛ばすたびに毎回ホストとセッションをつなぐ必要がなくなるので、めちゃ便利です。
import requests from requests.adapters import HTTPAdapter from urllib3.util.retry import Retry class RequestSession(): _has_instance = None def __new__(cls): if not cls._has_instance: cls._has_instance = super(RequestSession, cls).__new__(cls) cls.session = cls.create_session() return cls._has_instance @staticmethod def create_session(): session = requests.Session() retries = Retry(total=1) # このへんは別になくてもOK session.mount("http://", HTTPAdapter(max_retries=retries)) return session
では次にこれがちゃんとシングルトンになっているのかを確認します。
>>> rs = RequestSession() >>> id(rs) 140271749279264 >>> rs2 = RequestSession() >>> id(rs2) 140271749279264
識別子の値が一緒なので、同じインスタンスだとわかります。
ネットワークを確認する
ではこれを使用して、ネットワークを確認します。
しかし、Mac上にssコマンドやwatchコマンドがないので、Dockerコンテナ上で動かします。
FROM python:3 RUN apt-get -y install iproute2 watch RUN pip install --upgrade pip RUN pip install requests
次に、下記コマンドを叩きます。これで TCP の情報をリアルタイムで確認することが出来ます。
docker container run -it <IMAGE ID> bash watch ss -aot
別窓でインタラクティブモードでpythonを起動します。
docker ps docker exec -it <CONTAINER ID> bash python
TCP の状態を確認するために、まずは通常の get を使います
import requests requests.get("https://google.com") requests.get("https://google.com") requests.get("https://yahoo.co.jp")
別窓のssコマンドにはこのように表示されたかと思います。
State Recv-Q Send-Q Local Address:Port Peer Address:Port TIME-WAIT 0 0 172.17.0.2:..... 172.217.25.196:https timer:(timewait,49sec,0) TIME-WAIT 0 0 172.17.0.2:..... 216.58.220.142:https timer:(timewait,49sec,0) TIME-WAIT 0 0 172.17.0.2:..... 216.58.220.142:https timer:(timewait,50sec,0) TIME-WAIT 0 0 172.17.0.2:..... 216.58.220.142:https timer:(timewait,50sec,0) TIME-WAIT 0 0 172.17.0.2:..... 172.217.25.196:https timer:(timewait,50sec,0) TIME-WAIT 0 0 172.17.0.2:..... 216.58.220.142:https timer:(timewait,49sec,0) TIME-WAIT 0 0 172.17.0.2:..... 172.217.25.196:https timer:(timewait,49sec,0) TIME-WAIT 0 0 172.17.0.2:..... 172.217.25.196:https timer:(timewait,50sec,0)
次はシングルトンで書いたほうでgetしてみます。
rs = RequestSession() rs.session.get("https://google.com") rs.session.get("https://google.com") rs.session.get("https://yahoo.co.jp")
State Recv-Q Send-Q Local Address:Port Peer Address:Port ESTAB 0 0 172.17.0.2:..... 172.217.25.196:https ESTAB 0 0 172.17.0.2:..... 216.58.220.142:https
requests.get()とsession.get()で、TCPのコネクション数がかなり違うかと思います。
シングルトンで実装することで何が嬉しいのか ()
ローカルから外部サーバーに通信するとき、ルーターのNAPT機能を使用しています。
ただ、リクエストをしすぎてポート番号が枯渇してしまった場合、色々と問題が起きてしまいます。
そのため、ホストとのセッションを維持し続けて通信することで、上記問題は多少増しになるよね!といった感じです。
(上記で説明したように、requests.get()などは毎回セッション繋ぎに行ってます)
このあたりをより詳しく知りたい方はKeep Aliveで調べると面白いと思います。
おまけ
2020年も最後ですね。最後になにかやっときたいな〜〜〜〜けどめんどくさいな〜〜〜と思ってたら社内の人からこれネタにできるやんと言われたので書いた感じです😇
TwitterのTLは今年買ってよかったものとか2020年の振り返りばっかで溢れてますが、逆張り陰キャなのでそういうのは書きません()
今年も私のブログをご覧いただきありがとうございました。来年もよろしくお願いします🥲 みなさん良いお年を!
GASを使ってSlackにスプレッドシートのグラフを送信する
Incoming WebhookのAbout的なことはここに書かれてます。 api.slack.com
前提
Incoming Webhook導入&設定済み
全体的なコード
var WEBHOOK_URL = ""; var DRIVE_FOLDER_ID = ""; function sender() { var sheet = SpreadsheetApp.getActiveSpreadsheet().getSheetByName(""); // テキスト送信 var values = get_values(sheet); var text = generate_text(values[0], values[1], values[2]); var payload = JSON.stringify({ text: text }); post(payload); // 画像送信 var folder = DriveApp.getFolderById(DRIVE_FOLDER_ID); save_charts(sheet, folder); var image_ids = get_image_ids(folder); for (var i = 0; i < image_ids.length; i++) { payload = JSON.stringify({ blocks: [ { type: "image", image_url: "https://drive.google.com/uc?id=" + image_ids[i], alt_text: "chart_image", }, ], }); post(payload); } } function get_values(sheet) { // スプシからフィールドを取得。ここのコードの中身は気にしないでOK var l = sheet.getRange("").getDisplayValues(); today = Utilities.formatDate(new Date(), "JST", "YYYY/M/d"); for (var i = 0; i < l.length; i++) { if (l[i][0].match(today)) { return [l[i][0], l[i][1], l[i][3]]; } } } function generate_text(v1, v2, v3) { // メッセージ作成 var text = Utilities.formatString("%s, %s, %s", v1, v2, v3); return text; } function save_charts(sheet, folder) { // グラフ画像をGoogle Driveに保存 var charts = sheet.getCharts(); for (var i = 0; i < charts.length; i++) { var chart_image = charts[i] .getBlob() .getAs("image/png") .setName("chart" + i + ".png"); file = folder.createFile(chart_image); // Slackで画像を送信するにはこの権限でないとエラーが起きる file.setSharing(DriveApp.Access.ANYONE_WITH_LINK, DriveApp.Permission.VIEW); } } function get_image_ids(folder) { // Google Driveから画像のURLを取得 var files = folder.getFiles(); var ids = new Array(); while (files.hasNext()) { var file = files.next(); ids.push(file.getId()); } return ids; } function post(payload) { // Slackに画像を投稿する var req = { method: "post", contentType: "application/json", payload: payload, }; UrlFetchApp.fetch(WEBHOOK_URL, req); } function remover() { // imageのアクセス権限がリンクを知っている人なら誰でも見ることができるようになっている。 // それがセキュリティ的によろしくないので、送信後画像をディレクトリから削除する var folder = DriveApp.getFolderById(DRIVE_FOLDER_ID); var image_ids = get_image_ids(folder); remove_files(image_ids); } function remove_files(image_ids) { // fileを削除 for (var i = 0; i < image_ids.length; i++) { var file = DriveApp.getFileById(image_ids[i]); file.setSharing(DriveApp.Access.PRIVATE, DriveApp.Permission.VIEW); file.setTrashed(true); } }
Incoming Webhookの補足
画像の送信はこのようなJSONを送信する必要がある。
{ method: "post", contentType: "application/json", payload: { blocks: [ { type: "image", image_url: "https://drive.google.com/uc?id=" + image_ids[i], alt_text: "chart_image", }, ], }, }
いつからか仕様が変わって、Blockというものになったらしい。画像を送信したいときは、以下のドキュメントに沿う必要がある。 api.slack.com
APIのテストはBlock Kit Builderからすることができる。
ここにJSONを貼り付けて、送信できるかどうか試してみるのもいいと思う。
それか、curlを使うのも全然ありではある。
pyKNPで構文解析し、結果を有向グラフにする
ほぼほぼこいつと同じ内容です。
ssabcire.hatenablog.com
import networkx as nx from pyknp import KNP def tag(text: str) -> (list, list): ''' return tag_ids: [(子基本句ID, 親基本句ID), ...] ''' knp = KNP() tag_list = knp.parse(text).tag_list() tag_ids = list() for tag in tag_list: # 各基本句へのアクセス if tag.parent_id != -1: tag_ids.append((tag.tag_id, tag.parent_id)) return tag_list, tag_ids def graph_to_image(tag_list, tag_ids, dg: nx.DiGraph): for u, v in tag_ids: dg.add_nodes_from([tag_list[u].midasi, tag_list[v].midasi]) dg.add_edge(tag_list[u].midasi, tag_list[v].midasi) nx.nx_agraph.to_agraph(dg).draw('graph.png', prog='dot') if __name__ == '__main__': text = "私はスポーツが好きだが、バスケは嫌い。でも彼はバスケがうまい" tag_list, tag_ids = tag(text) dg = nx.DiGraph() graph_to_image(tag_list, tag_ids, dg)
出力結果

KNPの出力結果と比較して、正しく構文木になっていることが確認できます。

Juman++を使って形態素解析を行う
Juman++を使って形態素解析を行い、頻出回数順にソートします
ちなみにTwitterのjsonを読み込んでます
import json import re from glob import glob from pyknp import Juman def counter(text, d): jumanapp = Juman() result = jumanapp.analysis(text) for mrph in result.mrph_list(): if mrph.genkei in d: d[mrph.genkei] = d[mrph.genkei] + 1 else: d[mrph.genkei] = 1 filenames = glob('/Users/ssab/go/src/research/twitter/json/*.json') d = dict() for i, filename in enumerate(filenames): f = open(filename, 'r') tweet_text = json.load(f)['full_text'] text = re.sub( r'(https?://[\w/:%#\$&\?\(\)~\.=\+\-]+)|(RT@.*?:)|([ | ])', '', tweet_text ) counter(text, d) f.close() print(sorted(d.items(), key=lambda x: x[1], reverse=True))
出力
[('する', 3), ('を', 3), ('に', 2), ('インタフェース', 2), ('若い', 1), ('うち', 1), ('から', 1), ('毎日', 1), ('コツコツ', 1), ('続ける', 1), ('こと', 1), ('で', 1), ('老後', 1), ('2000万', 1), ('円', 1), ('用意', 1), ('ぬ', 1), ('済む', 1), ('ようだ', 1), ('なる', 1), ('ソリューション', 1), ('同じだ', 1), ('メソッド', 1), ('シグネチャ', 1), ('持つ', 1), ('複数', 1), ('の', 1), ('embeded', 1), ('作れる', 1), ('様', 1), ('proposal', 1), ('どま', 1), ('そ', 1), ('は', 1), ('こうして', 1), ('生まれる', 1), ('。', 1)]