Azure Stream AnalyticsでEvent HubsがMalformed eventになる原因と解決策(C# JSONロケール問題)

Azure Stream Analytics で Event Hubs を入力にして可視化しようとしたら「Malformed event, incorrect serialization format」が出て止まってしまう…。このエラーは受信側の設定ミスに見えて、実は送信側で作った JSON がロケールの影響で壊れているケースが非常に多いです。原因の見つけ方から、最短で直す実装例までまとめます。

目次

起きている現象:「Malformed event, incorrect serialization format」とは

Azure Stream Analytics(以下 ASA)は、Event Hubs から受け取ったイベントを「入力のシリアル化形式(JSON/Avro/CSV など)」に従ってデシリアライズします。ここで受信したバイト列が想定の形式になっていないと、入力イベントを解釈できず次のようなエラーになります。

  • Malformed event
  • incorrect serialization format

今回のケースでは、C# のコンソールアプリから vehicle_id / timestamp / latitude / longitude を JSON として Event Hubs に送信し、ASA →(最終的に)Power BI に流したいのに、ASA 側で入力が「壊れたイベント」と判定されて止まっています。

結論:JSON は“正しそう”に見えても、ロケールで簡単に壊れる

問題の本質は、送信側が JSON を文字列連結で手書きしている点です。質問にある抜粋コードは次のような形でした。

var latitude = random.NextDouble() * 180 - 90;
var longitude = random.NextDouble() * 360 - 180;

var data = "{\"vehicle_id\": " + vehicleId + ", \"timestamp\": \"" + timestamp + "\", \"latitude\": " + latitude + ", \"longitude\": " + longitude + "}";

このコードは一見正しく見えますが、実行環境のロケール(地域設定)によっては、double を文字列化したときの小数点が「.(ドット)」ではなく「,(カンマ)」になります。その結果、送信される JSON が次のように壊れます。

{
  "vehicle_id": 7941,
  "timestamp": "2023-10-19 19:25:17",
  "latitude": -45,8807284588388,
  "longitude": 90,01136751036734
}

JSON の数値表現では 小数点はドットのみが許容されます。上の -45,8807... は「-45 と 8807... がカンマ区切りで並んでいる」と解釈され、JSON として不正になります。ASA は入力を JSON として解析できず、結果として 「Malformed event, incorrect serialization format」になります。

なぜこうなるのか:double の文字列化はカルチャ(Culture)に依存する

C# では、数値を文字列にするときに CurrentCulture(現在のカルチャ)のルールが適用されます。たとえば一部の地域設定(例:ドイツ語圏、フランス語圏など)では、

  • 小数点:,(カンマ)
  • 桁区切り:.(ドット) やスペース

が一般的です。つまり latitude.ToString() 相当の処理が走ると、JSON の仕様と衝突しやすい形式で文字列化されます。

「自分の PC では再現しない」のに、本番サーバーや別の実行環境で突然発生するのもこのパターンです。コンテナ、VM、ホスト OS の言語設定、実行ユーザーのカルチャが変わるだけで、出力が変わります。

最初にやるべき切り分け:送信した“生のJSON”を必ず確認する

ASA 側の設定を疑う前に、まずは「Event Hubs に実際に送られているバイト列が、正しい JSON なのか」を確認します。ここが曖昧なまま設定をいじると、遠回りになりがちです。

確認ポイント狙い具体的なやり方
送信前に data をログ出力アプリが作った文字列が壊れていないかConsole.WriteLine(data) でまず確認
JSON パーサで検証“見た目OK”でもパースできるかJToken.Parse(data) や JsonDocument.Parse(...) を試す
別の受信アプリで Event Hubs から読むEvent Hubs に入っている実データの確認簡単なコンシューマ(後述)で受信し、文字列として表示
ASA の入力メトリクス確認入力デシリアライズで落ちているかジョブの Monitoring で “Deserialization errors” の増加を確認

再現テスト:わざとロケールを変えて壊れるか確認する

原因がカルチャ依存かどうかは、送信アプリの冒頭でカルチャを変えると明確になります(検証用です。恒久対応は後述の方法を推奨します)。

using System.Globalization;
using System.Threading;

// 例:小数点がカンマになりやすいカルチャ
Thread.CurrentThread.CurrentCulture = new CultureInfo("fr-FR");
Thread.CurrentThread.CurrentUICulture = new CultureInfo("fr-FR");

この状態で元の「文字列連結 JSON」を使って送信すると、latitude/longitude がカンマ区切りになり、JSON が壊れやすくなります。

最も確実な解決策:JSON を手書きせずシリアライザで生成する

結論としてはこれがベストです。

  • 文字列連結で JSON を作らない
  • C# のオブジェクトを作り、JSON シリアライザで文字列化する

シリアライザは JSON 仕様に沿ってエスケープや数値表現を行うため、ロケールに依存せず常に正しい JSON を生成できます。

Newtonsoft.Json(Json.NET)での修正版:そのまま置き換えやすい

既存のサンプルに近い形で、元の var data = "{..."; を置き換える例です。snake_case のキー(vehicle_id)を維持したい場合は JsonProperty を付けると安全です。

using System;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using Azure.Messaging.EventHubs;
using Azure.Messaging.EventHubs.Producer;
using Newtonsoft.Json;

class Program
{
    private const string connectionString = "[secret]";
    private const string eventHubName = "[secret]";

    static async Task Main(string[] args)
    {
        await using var producerClient = new EventHubProducerClient(connectionString, eventHubName);
        var random = new Random();

        while (true)
        {
            var message = new VehicleMessage
            {
                VehicleId = random.Next(1000, 10000),
                Timestamp = DateTime.UtcNow,               // DateTime のままでOK
                Latitude = random.NextDouble() * 180 - 90,
                Longitude = random.NextDouble() * 360 - 180
            };

            // ロケール非依存で正しい JSON が生成される
            var json = JsonConvert.SerializeObject(message);

            using var eventBatch = await producerClient.CreateBatchAsync();

            // 可能なら ContentType を入れると、後から調査が楽になります(必須ではありません)
            var eventData = new EventData(Encoding.UTF8.GetBytes(json));
            eventData.ContentType = "application/json";

            if (!eventBatch.TryAdd(eventData))
            {
                Console.WriteLine("Event is too large for the batch.");
                continue;
            }

            await producerClient.SendAsync(eventBatch);
            Console.WriteLine($"Sent: {json}");

            Thread.Sleep(2000);
        }
    }
}

public class VehicleMessage
{
    [JsonProperty("vehicle_id")]
    public int VehicleId { get; set; }

    [JsonProperty("timestamp")]
    public DateTime Timestamp { get; set; }

    [JsonProperty("latitude")]
    public double Latitude { get; set; }

    [JsonProperty("longitude")]
    public double Longitude { get; set; }
}

ポイントは次のとおりです。

  • double の小数点がカルチャで変わっても、シリアライザが JSON 仕様の表現で出力する
  • DateTime を文字列連結で埋め込まず、オブジェクトのまま持たせる(ISO 8601 形式で出力されやすく、解析もしやすい)
  • Encoding.UTF8.GetBytes を使い、UTF-8 として送る

System.Text.Json 版:追加ライブラリなしで実装したい場合

.NET の標準ライブラリで完結させたい場合は System.Text.Json も有力です。キー名を合わせたい場合は JsonPropertyName を使います。

using System;
using System.Text;
using System.Text.Json;
using System.Text.Json.Serialization;
using System.Threading;
using System.Threading.Tasks;
using Azure.Messaging.EventHubs;
using Azure.Messaging.EventHubs.Producer;

class Program
{
    private const string connectionString = "[secret]";
    private const string eventHubName = "[secret]";

    static async Task Main(string[] args)
    {
        await using var producerClient = new EventHubProducerClient(connectionString, eventHubName);
        var random = new Random();

        var jsonOptions = new JsonSerializerOptions
        {
            // 必要なら読みやすい整形も可能(本番は false 推奨)
            WriteIndented = false
        };

        while (true)
        {
            var message = new VehicleMessage
            {
                VehicleId = random.Next(1000, 10000),
                Timestamp = DateTime.UtcNow,
                Latitude = random.NextDouble() * 180 - 90,
                Longitude = random.NextDouble() * 360 - 180
            };

            var json = JsonSerializer.Serialize(message, jsonOptions);

            using var eventBatch = await producerClient.CreateBatchAsync();
            var eventData = new EventData(Encoding.UTF8.GetBytes(json))
            {
                ContentType = "application/json"
            };

            eventBatch.TryAdd(eventData);
            await producerClient.SendAsync(eventBatch);

            Console.WriteLine($"Sent: {json}");
            Thread.Sleep(2000);
        }
    }
}

public class VehicleMessage
{
    [JsonPropertyName("vehicle_id")]
    public int VehicleId { get; set; }

    [JsonPropertyName("timestamp")]
    public DateTime Timestamp { get; set; }

    [JsonPropertyName("latitude")]
    public double Latitude { get; set; }

    [JsonPropertyName("longitude")]
    public double Longitude { get; set; }
}

対策の選び方(おすすめ順)

対策おすすめ度メリット注意点
シリアライザで JSON 生成(Newtonsoft / System.Text.Json)最優先ロケール・エスケープ・数値表現をまとめて安全化できるキー名の互換が必要なら属性で調整
InvariantCulture を使って文字列連結を延命次点既存の実装を最小差分で直せる将来別の項目追加でまた壊れやすい(エスケープ漏れ等)
ASA 側の設定変更で吸収非推奨送信側を触れない事情がある場合に検討壊れた JSON は設定変更で正しくはならない

(参考)どうしても手書きするなら:CultureInfo.InvariantCulture を固定する

事情があって「まずは最小改修で直したい」という場合、文字列化の時点で InvariantCulture を使えば、小数点はドットになります。ただし、JSON の手書きは他のミス(クォート、エスケープ、null、配列、入れ子)が増えやすいので、恒久対応としてはシリアライザが安全です。

using System.Globalization;

var latitude = random.NextDouble() * 180 - 90;
var longitude = random.NextDouble() * 360 - 180;

var data = string.Format(
    CultureInfo.InvariantCulture,
    "{{\"vehicle_id\": {0}, \"timestamp\": \"{1}\", \"latitude\": {2}, \"longitude\": {3}}}",
    vehicleId,
    timestamp,
    latitude,
    longitude
);

ここで重要なのは、小数点が必ず “.” になることです。逆に言えば、元の実装は「たまたま自分の環境で “.” だった」だけで、環境が変わると壊れる可能性を常に抱えていました。

修正後も ASA のエラーが残る場合:古い“壊れたイベント”が残っている

送信側を直しても、ASA のエラーがすぐ消えないことがあります。よくある理由は次のとおりです。

Event Hubs の保持期間(Retention)により、過去の不正メッセージが読み続けられる

Event Hubs は設定された保持期間の間、イベントを保持します。修正前に送った壊れた JSON が残っていると、ASA はそれらも順に読み込み、入力デシリアライズで失敗し続けます。

この場合の対処は次の方向性になります。

  • 不正データが通過するまで待つ(送信量が少ない検証環境ならこれが最も簡単)
  • ASA の入力の開始位置を「現在」からにする(再起動や入力設定の見直し)
  • 別のコンシューマグループに切り替えて、クリーンな位置から読み直す

チェックポイントが原因で、同じ不正イベントを何度も読む

ASA は処理位置を管理します。失敗が続いたタイミングによっては、同じ地点からリトライし続ける挙動に見えることがあります。入力の開始位置・コンシューマグループ・ジョブ再起動の組み合わせで切り分けると、原因が見えやすくなります。

Stream Analytics 側の設定も一応確認:ただし“壊れた JSON”は救えない

今回の根本原因は送信側ですが、同時に「ASA 側が想定通り JSON として受け取る設定になっているか」も押さえておくと、次回以降のトラブルシュートが速くなります。

項目推奨意図
Input の SerializationJSON送信が JSON の場合は必須
EncodingUTF-8送信側も UTF-8 に統一する
Compressionなし(使うなら送受で一致)不一致だと読めない
Timestamp の扱いUTC / ISO 8601 を推奨時刻の解釈ズレを防ぐ
Consumer Group検証用と本番用を分けるデバッグ時に読み位置が干渉しにくい

なお、入力が壊れているときに ASA 側で起きるのは「形式エラー」なので、SQL クエリや Power BI 出力設定の前段で止まります。まずは “正しい JSON を Event Hubs に入れる” を最優先してください。

デバッグを一気に楽にする:Event Hubs 受信用の最小コンシューマ

「ASA が何を読んで落ちているか」を可視化するには、Event Hubs から直接受信して中身を出すのが最短です。下の例は “届いた JSON をそのまま表示する” だけに絞ったイメージです(運用用途ではチェックポイントや例外処理を強化してください)。

using System;
using System.Text;
using System.Threading.Tasks;
using Azure.Messaging.EventHubs.Consumer;

class Receiver
{
    static async Task Main()
    {
        var connectionString = "[secret]";
        var eventHubName = "[secret]";
        var consumerGroup = EventHubConsumerClient.DefaultConsumerGroupName;

        await using var consumer = new EventHubConsumerClient(consumerGroup, connectionString, eventHubName);

        await foreach (var partitionEvent in consumer.ReadEventsAsync())
        {
            var body = Encoding.UTF8.GetString(partitionEvent.Data.EventBody.ToArray());
            Console.WriteLine(body);

            // ここで JSON パースを試して、壊れているイベントを即特定できる
            // 例:Newtonsoft.Json.Linq.JToken.Parse(body);
        }
    }
}

この受信側で、latitude や longitude が -45,8807... のように出てきたら、ASA ではなく送信側が原因だと確定できます。

“文字列連結 JSON”が危険な理由:今回以外にも地雷が多い

今回のロケール問題は典型例ですが、手書き JSON には他にも落とし穴があります。今後フィールドが増えるほど事故率が上がるので、「なぜシリアライザ推奨なのか」を整理しておきます。

落とし穴例起きる問題
ロケール依存小数点がカンマになるJSON が壊れて ASA でデシリアライズ不能
エスケープ漏れ文字列にダブルクォートが入るJSON 構造が崩壊
null/空文字の扱いフィールド欠落や余計なカンマパースできない/型推論がおかしくなる
日時形式がバラバラ環境で ToString() の形式が変わるASA/Power BI 側で日時として扱えない
浮動小数の特殊値NaN/InfinityJSON 仕様外になり得る(処理系によって挙動が違う)

シリアライザはこれらの多くを標準で安全に扱えます。特に IoT・テレメトリ系のパイプライン(Event Hubs → ASA → Power BI)はイベント数が増えやすく、一度壊れたデータが混ざると調査時間が膨らみます。最初から安全な生成方法に寄せるのが結果的に早いです。

Power BI 可視化まで見据えた、データ設計の実務ポイント

今回のテーマは「ASA の Malformed event」ですが、同じプロジェクトで詰まりやすいポイントも合わせて押さえておくと、実装が安定します。

timestamp は UTC を基本にし、ISO 8601 を意識する

送信データに時刻を持たせる場合、UTC に統一するのが安全です。シリアライザを使えば、DateTime.UtcNow は ISO 8601(例:2025-12-16T12:34:56.789Z)に近い形式で出力されることが多く、解析・変換が容易です。

フィールド名と型は“固定”し、途中で変えない

Power BI 側は型推論や列の扱いでハマることがあります。運用を考えるなら、少なくとも次を安定させると良いです。

  • vehicle_id:数値(int)
  • timestamp:UTC の日時(DateTime/DateTimeOffset 由来)
  • latitude/longitude:数値(double)

“イベントのバージョン”を入れると将来の変更が楽になる

実務では、後からフィールドが増えたり、名前が変わったりします。小さくてもよいので、次のようなフィールドを入れておくと、ASA/Power BI だけでなく、将来別の下流(Data Lake、Kusto、Synapse など)に流したときにも整合が取りやすくなります。

  • schema_version(例:1)
  • source(例:simulator / device)

運用で再発させないためのチェック:テストと監視

「直ったはずなのに、また Malformed event が出た」を防ぐには、送信側の段階で弾く仕組みが効果的です。

送信前に JSON としてパースできるかを軽く検証する

大量送信の本番パスで毎回パースするとコストが増えますが、検証環境・ステージング・サンプル率を決めた監視などで “壊れた JSON を出していないか” を見張れます。

// 例:Newtonsoft.Json を使った軽い検証(検証用途)
using Newtonsoft.Json.Linq;

try
{
    JToken.Parse(json); // パースできなければ例外
}
catch (Exception ex)
{
    Console.WriteLine($"Invalid JSON: {ex.Message}");
}

ASA のメトリクス(デシリアライズエラー)をアラートにする

ASA 側で “入力デシリアライズエラーが増えたら通知” を作っておくと、送信側の変更や環境差分による事故を早期に検知できます。エラーが出てから Power BI の画面で気付くより遥かに早く手当てできます。

まとめ:原因は「ASA ではなく、カルチャで壊れた JSON」

  • 「Malformed event, incorrect serialization format」は、ASA が入力を想定形式(今回なら JSON)として解釈できないときに発生する
  • 文字列連結で JSON を作ると、double の小数点がロケールの影響でカンマになり、JSON が壊れる
  • 最も確実な対策は、Json.NET / System.Text.Json などのシリアライザで JSON を生成する
  • 修正後もエラーが残る場合は、Event Hubs の保持期間内にある “過去の不正イベント” を ASA が読んでいる可能性を疑う

Event Hubs → ASA → Power BI の流れは、入口(送信側)の JSON がすべての前提になります。まずは「送信している JSON が仕様通りか」を確実にし、そのうえで ASA の入力設定・開始位置・可視化を整えると、最短で安定運用に乗せられます。

この記事を書いた人

実務の現場で詰まりがちなポイントを地図にするITブログ「IT trip」を運営。Windows/Office(Teams・Excel)からSQL、サーバ運用、ガジェットまで、再現性のある手順と“なぜそうなるか”を丁寧に解説します。読んだらすぐ試せること、そして迷った人の次の一歩が見えることを大切にしています。

コメント

コメントする

目次