たのしい駆動開発

たのしいアウトプットの場所

Cloud Data Fusionの機能でつまづいた部分について

Google Cloud Platform Advent Calendar 2021の22日目カレンダー2の記事です。


最近Cloud Data Fusionをはじめて触りました。GUIベースで結構感覚的に処理ができるので、結構いいツールだと思います。
そして今回はCloud Data Fusionで説明が必要だなあと感じた機能について説明します。
Cloud Data Fusionについての説明や、インスタンスの作成などは公式ドキュメントを参考にしてください。



今回使用するサンプルのテーブルとデータは下記を使用します。

f:id:ssabcire:20211222093656p:plain

RawDenormalizer

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

Wranglerでの列Joinはこのように行います
f:id:ssabcire:20211222093912p:plain f:id:ssabcire:20211222093924p:plain

そして、RowDenormalizerでこのように指定します。
f:id:ssabcire:20211222094130p:plain

RowDenormalizerをした結果はこのようなデータになります。 f:id:ssabcire:20211222095727p:plain

Join

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

f:id:ssabcire:20211222094743p:plain

テーブル出力

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

f:id:ssabcire:20211222095206p:plain

ちなみに、デフォルトではパイプラインを実行するたびにテーブルに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)



出力結果 f:id:ssabcire:20191022222905p:plain



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

Juman++を使って形態素解析を行う

Juman++を使って形態素解析を行い、頻出回数順にソートします
ちなみにTwitterjsonを読み込んでます

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)]