Avroで壊れないイベントschemaを設計する|writer・reader・互換性・Python

「Avroで壊れないイベントschemaを設計する|writer・reader・互換性・Python」の内容を表す技術イラスト

システム間でeventやrecordを長期間受け渡すと、送信側と受信側を同じ日に更新できないことがあります。新しいfieldを追加しただけでも、古いconsumerが読めなくなれば連携は停止します。

Apache Avroは、JSONで定義したschemaとbinary dataを組み合わせるserialization方式です。dataを書いたときのwriter schemaと、読む側が期待するreader schemaを照合することで、一定範囲のschema変更を吸収できます。この記事では、型付きeventの構造、互換性規則、Python実装、再処理とversion管理を整理します。

目次

Avroで解決すること

JSON単体では、1が32bit整数か64bit整数か、日時文字列がどの精度か、field追加をどう扱うかは別の契約が必要です。Avroではrecord、field、type、defaultをschemaに定義します。

Apache Avro 1.12.0仕様が定めるprimitive typeは、null、boolean、int、long、float、double、bytes、stringです。record、enum、array、map、union、fixedも組み合わせられます。

Avro Object Container Fileでは、file headerにwriter schemaが保存され、recordはbinary encodingされたblockとして格納されます。message brokerで1 recordずつ送る場合は、schemaそのものではなくschema IDやfingerprintをpayloadへ関連付ける運用が必要です。

最小の測定event schema

圧力測定eventのversion 1を次のschemaにします。

{
  "type": "record",
  "name": "Measurement",
  "namespace": "com.example.engineering",
  "fields": [
    {"name": "measurement_id", "type": "string"},
    {
      "name": "measured_at",
      "type": {"type": "long", "logicalType": "timestamp-micros"}
    },
    {"name": "pressure_pa", "type": "long"},
    {"name": "valid", "type": "boolean"},
    {"name": "note", "type": ["null", "string"], "default": null}
  ]
}

schema全体がrecord、fieldsがfield定義のarrayです。各fieldはnameとtypeを持ちます。noteはnullまたはstringを許すunionです。

最小の入力recordは次の形です。

{
  "measurement_id": "m-000184",
  "measured_at": "2026-10-09T01:30:00Z",
  "pressure_pa": 14000000,
  "valid": true,
  "note": null
}

JSONは入力例であり、Avro binary内の日時は文字列ではありません。timestamp-microsはUnix epochからのmicrosecond数を持つlongへ意味を付けるlogical typeです。入力境界でtimezone付き日時をUTCへ正規化します。

writer schemaとreader schemaを分ける

writer schemaはdataをserializeしたときのschema、reader schemaは現在のapplicationが期待するschemaです。両者を同じものとして固定すると、過去dataを読むたびに旧applicationが必要になります。

version 2では校正記録への参照を追加します。

{
  "type": "record",
  "name": "Measurement",
  "namespace": "com.example.engineering",
  "fields": [
    {"name": "measurement_id", "type": "string"},
    {
      "name": "measured_at",
      "type": {"type": "long", "logicalType": "timestamp-micros"}
    },
    {"name": "pressure_pa", "type": "long"},
    {"name": "valid", "type": "boolean"},
    {"name": "note", "type": ["null", "string"], "default": null},
    {
      "name": "calibration_id",
      "type": ["null", "string"],
      "default": null
    }
  ]
}

旧dataにはcalibration_idがありませんが、reader schemaにdefaultがあるためnullとして補えます。defaultは「書き込み時にfieldを省略してよい」という意味ではなく、schema解決時にwriter側へ存在しないfieldを補う値です。

schema解決の基本規則

Avro仕様ではrecord fieldを宣言順ではなくnameで対応付けます。

変更readerの動作注意点
default付きfieldを追加defaultで補完defaultの意味を契約化する
writerだけにあるfieldreaderは無視dataを再保存すると失われ得る
defaultなしfieldをreaderへ追加読取error既存dataと非互換
field順を変更nameで照合record名とfield名を維持する
intからlongへ拡張promotion可能逆方向は保証されない
field名を変更原則別field必要ならreader側aliasを検討

["null", "string"]のようなunionでは、どのbranchへ値を対応させるかもschemaの一部です。defaultの型がunion内の型と一致することを確認し、optionalのつもりならNULLとfield未指定の意味も分けます。

Pythonで旧schemaを書き新schemaで読む

次の例はversion 1でfileを書き、calibration_idを追加したversion 2で読みます。Python 3では公式のavro packageを使います。

import json
from datetime import datetime, timezone

import avro.schema
from avro.datafile import DataFileReader, DataFileWriter
from avro.io import DatumReader, DatumWriter


V1 = {
    "type": "record",
    "name": "Measurement",
    "namespace": "com.example.engineering",
    "fields": [
        {"name": "measurement_id", "type": "string"},
        {
            "name": "measured_at",
            "type": {"type": "long", "logicalType": "timestamp-micros"},
        },
        {"name": "pressure_pa", "type": "long"},
        {"name": "valid", "type": "boolean"},
        {"name": "note", "type": ["null", "string"], "default": None},
    ],
}

V2 = {
    **V1,
    "fields": V1["fields"]
    + [
        {
            "name": "calibration_id",
            "type": ["null", "string"],
            "default": None,
        }
    ],
}

writer_schema = avro.schema.parse(json.dumps(V1))
reader_schema = avro.schema.parse(json.dumps(V2))

with open("measurements.avro", "wb") as stream:
    with DataFileWriter(stream, DatumWriter(), writer_schema) as writer:
        writer.append(
            {
                "measurement_id": "m-000184",
                "measured_at": datetime(
                    2026, 10, 9, 1, 30, tzinfo=timezone.utc
                ),
                "pressure_pa": 14_000_000,
                "valid": True,
                "note": None,
            }
        )

with open("measurements.avro", "rb") as stream:
    datum_reader = DatumReader(readers_schema=reader_schema)
    with DataFileReader(stream, datum_reader) as reader:
        records = list(reader)

assert records[0]["calibration_id"] is None

Apache AvroのPython入門資料でも、DataFileWriter、DatumWriter、DataFileReader、DatumReaderを使う基本形が示されています。fileはbinary modeで開きます。

正常系と異常系をテストする

schema fileがJSONとしてparseできるだけでは、互換性を確認したことになりません。最低限、次を自動テストします。

  • version 1で書いたrecordをversion 2で読める
  • 追加fieldへ期待したdefaultが入る
  • 必須field欠落、型違い、未知のenum symbolを拒否する
  • falseと0、NULLと空文字を混同しない
  • timezoneなし日時を入力境界で拒否する
  • measurement_id重複をDBまたは受信処理で検出する
  • 未対応schema versionを隔離する

Avroの型検証は、圧力が業務上の許容範囲か、IDが一意か、単位がPaかまでは判断しません。構造validationと業務validationを分けます。

logical typeと単位の落とし穴

logical typeはprimitive typeへ意味を追加します。日時ならtimestamp-micros、十進数ならdecimalを利用できます。ただし、利用する言語の実装がそのlogical typeをどのobjectへ変換するかを確認します。

Decimalではprecisionとscaleも互換性判定へ影響します。floatへ変換してから保存せず、十進数として扱います。圧力のような単位はAvro標準のlogical typeではないため、pressure_paのようなfield名、schema property、またはdata contractで固定します。

schemaをIDと履歴で管理する

schema本文を毎messageへ添付するとpayloadが大きくなります。運用では次を管理します。

  • schemaの不変IDまたはfingerprint
  • schema本文と登録日時
  • backward・forwardの互換性方針
  • producerとconsumerが利用中のversion
  • 廃止予定fieldと移行期限
  • schema変更のtest結果

Object Container Fileにはwriter schemaが入ります。一方、single-object encodingはmarker、schema fingerprint、binary dataで構成されます。受信側がfingerprintから正しいwriter schemaを解決できる仕組みが必要です。

schema IDをpayloadの業務IDとして使ってはいけません。measurement_idはrecordの同一性、schema IDは構造の同一性を表します。

冪等性と再処理

Avroはserialization方式であり、重複排除を自動では行いません。messageの再送に備えて不変のevent IDを持たせ、受信側の一意制約で処理済みを判定します。

過去recordを新schemaへ変換するときは、元dataを上書きせず、writer schema ID、変換program version、出力schema IDを記録します。同じ入力と変換versionから同じ出力へ収束するようにし、途中失敗後もbatch IDから再開できる設計にします。

セキュリティと公開情報の分離

Avro binaryは暗号化ではありません。API key、token、passwordをrecordやschema metadataへ入れず、転送路とstorage側で暗号化・アクセス制御します。

schema自体にも内部system名、顧客名、非公開URLを埋め込まないようにします。信頼できないpayloadではmessage size、array要素数、string長、nested depthを制限し、解釈前に利用可能なschema IDか確認します。

CAD・BOM・APIへ再利用する

CAD属性更新、BOM revision、計算完了、測定結果を共通event envelopeで表し、業務ごとの差をrecord schemaへ閉じ込められます。大きなCAD file自体は埋め込まず、file ID、revision、content hashを渡して認可済みstorageから取得します。

Avro schemaからDB、API、CSVへ変換するときは、longの範囲、NULL、timestamp、Decimal、単位をmapping表として残します。形式を変えても意味が変わらないことが、再利用可能な設計dataの条件です。

まとめ

Apache Avroの中心は、binary化そのものではなく、writer schemaとreader schemaを使った型付きdata契約です。

  • recordのfield名、型、defaultをschemaへ明記する
  • 追加fieldには既存dataを読めるdefaultを設計する
  • defaultと書込時の省略可能性を混同しない
  • readerがwriter schemaを解決できるようIDやfingerprintを管理する
  • logical type、NULL、単位、IDの意味を別途検証する
  • 新旧schemaの組み合わせを自動テストする
  • event IDによる重複排除と、変換履歴による再処理を設計する
  • secretと機密metadataをdata contractから分離する

この仕組みをAPI、message broker、file、CAD、BOMへ共通化すれば、producerとconsumerを同時更新できない環境でも、過去dataを保持しながらschemaを段階的に変更できます。

参考情報

参考になったらシェアしてください
  • URLをコピーしました!
  • URLをコピーしました!

この記事を書いた人

機械設計・油圧・CAD・Python・AIなど、ものづくりに関わる技術を扱っています。工学知識を整理・構造化し、設計や自動化に再利用できる形へ変えていくことを目指しています。

目次