すべてのプロダクト
Search
ドキュメントセンター

Realtime Compute for Apache Flink:Log Service (SLS)

最終更新日:Aug 14, 2026

Log Service (SLS) コネクタの使用方法について説明します。

背景情報

Simple Log Service は、ログデータのエンドツーエンドサービスです。ログデータの収集、消費、転送、クエリ、分析を効率的に行うことができます。これにより、O&M 効率が向上し、大量のログデータを処理できるようになります。

次の表に、SLS コネクタの機能を示します。

カテゴリ

説明

サポートされるタイプ

ソーステーブルと結果テーブル

実行モード

ストリーミングモードのみ

コネクタ固有のメトリック

N/A

データフォーマット

N/A

API タイプ

SQL、DataStream API、およびデータインジェスト YAML API

結果テーブルでのデータの更新または削除

結果テーブルは追加専用です。データを更新または削除することはできません。

特徴

SLS ソースコネクタは、メッセージ属性フィールドを直接読み取ります。次の表に、サポートされているフィールドを示します。

パラメーター

タイプ

説明

__source__

STRING METADATA VIRTUAL

メッセージソース。

__topic__

STRING METADATA VIRTUAL

メッセージトピック。

__timestamp__

BIGINT METADATA VIRTUAL

ログ時間。

__tag__

MAP<VARCHAR, VARCHAR> METADATA VIRTUAL

メッセージタグ。

例えば、属性 "__tag__:__receive_time__":"1616742274" の場合、'__receive_time__' と '1616742274' はマップにキーと値のペアとして格納されます。SQL で値にアクセスするには、__tag__['__receive_time__'] を使用します。

前提条件

Log Service プロジェクトと Logstore が作成されていることを確認してください。詳細については、「プロジェクトと Logstore の作成」をご参照ください。

制限事項

  • YAML で定義されたデータインジェストの同期ソースとして SLS を使用することは、Ververica Runtime (VVR) 11.1 以降でのみサポートされています。

  • SLS コネクタは、at-least-once セマンティクスのみを保証します。

  • source parallelism をシャード数より高く設定しないでください。リソースの無駄になります。さらに、VVR 8.0.5 以前では、シャード数の変更により自動 failover 機能が失敗し、一部のシャードが消費されなくなる可能性があります。

SQL

構文

CREATE TABLE sls_table(
  a INT,
  b INT,
  c VARCHAR
) WITH (
  'connector' = 'sls',
  'endPoint' = '<yourEndPoint>',
  'project' = '<yourProjectName>',
  'logStore' = '<yourLogStoreName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}'
);

オプション付き

  • 全般

    パラメーター

    説明

    タイプ

    必須

    デフォルト

    備考

    connector

    使用するコネクタ。

    String

    はい

    なし

    これを sls に設定します。

    endPoint

    Log Service (SLS) のエンドポイント。

    String

    はい

    なし

    Log Service (SLS) の VPC アクセスアドレスを指定します。詳細については、「サービスエンドポイント」をご参照ください。

    説明
    • デフォルトでは、Realtime Compute for Apache Flink はインターネットにアクセスできません。ご利用の Virtual Private Cloud (VPC) からインターネットへのアクセスを有効にするには、NAT Gateway を使用します。詳細については、「インターネットへのアクセス方法」をご参照ください。

    • インターネット経由で SLS にアクセスすることは推奨されません。もしそうする必要がある場合は、HTTPS を使用し、転送アクセラレーションを有効にしてください。詳細については、「転送アクセラレーションの管理」をご参照ください。

    project

    SLS プロジェクトの名前。

    String

    はい

    なし

    N/A

    logStore

    Logstore または MetricStore の名前。

    String

    はい

    なし

    Logstore 内のデータは、MetricStore 内のデータと同じ方法で消費されます。

    accessId

    ご利用の Alibaba Cloud アカウントの AccessKey ID。

    String

    はい

    なし

    詳細については、「AccessKey ペアの取得」をご参照ください。

    重要

    AccessKey ペアが公開されるのを防ぐため、変数を使用して AccessKey ID と AccessKey Secret を指定することを推奨します。詳細については、「プロジェクト変数」をご参照ください。

    accessKey

    ご利用の Alibaba Cloud アカウントの AccessKey Secret。

    String

    はい

    なし

  • ソース固有

    パラメーター

    説明

    タイプ

    必須

    デフォルト

    備考

    enableNewSource

    FLIP-27 インターフェイスを実装した新しいデータソースを使用するかどうかを指定します。

    Boolean

    いいえ

    false

    新しいソースは、シャードの変更に自動的に適応し、すべてのソースサブタスクにシャードをできるだけ均等に分散させることができます。

    重要
    • このオプションは VVR 8.0.9 以降でのみサポートされています。VVR 11.1 以降のデフォルト値は true です。

    • このオプションの値を変更すると、保存された状態からジョブを復元できなくなります。履歴オフセットから消費を再開するには、まず consumerGroup オプションでジョブを開始して、SLS コンシューマーグループに消費の進行状況を記録します。次に、consumeFromCheckpoint オプションを true に設定し、状態なしでジョブを再起動します。

    • Logstore に読み取り専用のシャードがある場合、一部のサブタスクは自身のシャードを終えた後も他のシャードからデータを要求し続けることがあります。これにより、ワークロードが不均衡になり、パフォーマンスに影響を与える可能性があります。この問題を軽減するには、並列度を調整したり、スケジューリング戦略を最適化したり、小さなシャードをマージしてシャード数を減らし、タスクの割り当てを簡素化したりすることができます。

    shardDiscoveryIntervalMs

    動的なシャード検出の間隔。

    Long

    いいえ

    60000

    動的検出を無効にするには、このオプションに負の値を設定します。単位:ミリ秒。

    説明
    • 値は 60,000 ミリ秒 (1 分) 以上である必要があります。

    • このオプションは、enableNewSourcetrue に設定されている場合にのみ有効です。

    • このオプションは VVR 8.0.9 以降でのみサポートされています。

    startupMode

    ソーステーブルの起動モード。

    String

    いいえ

    timestamp

    • timestamp (デフォルト):指定された開始時刻からログの消費を開始します。

    • latest:最新のオフセットからログの消費を開始します。

    • earliest:最も古いオフセットからログの消費を開始します。

    • consumer_group:コンシューマーグループによって記録されたオフセットからログの消費を開始します。コンシューマーグループがシャードの消費オフセットを記録していない場合、消費は最も古いオフセットから開始されます。

    重要
    • VVR バージョン 11.1 より前では、consumer_group の値はサポートされていません。consumeFromCheckpointtrue に設定する必要があります。この場合、ログの消費は指定されたコンシューマーグループによって記録されたオフセットから開始され、起動モードの設定は有効になりません。

    startTime

    ログ消費の開始時刻。

    String

    いいえ

    現在時刻

    フォーマットは yyyy-MM-dd hh:mm:ss です。

    これは startupModetimestamp に設定されている場合にのみ有効です。

    説明

    startTime および stopTime オプションは、SLS の __timestamp__ 属性ではなく、__receive_time__ 属性に基づいています。

    stopTime

    ログ消費の終了時刻。

    String

    いいえ

    なし

    フォーマットは yyyy-MM-dd hh:mm:ss です。

    説明
    • このオプションは履歴ログの消費にのみ使用され、過去の時刻に設定する必要があります。未来の時刻に設定すると、新しいログが書き込まれない場合に消費が途中で停止し、エラーメッセージなしでデータフローが中断される可能性があります。

    • すべてのログが消費された後に Flink ジョブを終了させたい場合は、exitAfterFinishtrue に設定する必要があります。

    consumerGroup

    コンシューマーグループの名前。

    String

    いいえ

    なし

    コンシューマーグループは消費の進行状況を記録するために使用されます。固定フォーマットなしでカスタム名を指定できます。

    説明

    異なる Flink ジョブは異なるコンシューマーグループを使用する必要があります。複数の Flink ジョブが同じコンシューマーグループを使用する場合、それらは協調せず、各ジョブがすべてのデータを消費します。これは、Flink が SLS からデータを消費する際に、パーティション割り当てに SLS コンシューマーグループを使用しないためです。その結果、同じコンシューマーグループを共有していても、各コンシューマーは独立してメッセージを消費します。

    consumeFromCheckpoint

    コンシューマーグループのチェックポイントから消費するかどうかを指定します。

    String

    いいえ

    false

    • true:コンシューマーグループも指定する必要があります。Flink プログラムは、コンシューマーグループに保存されているチェックポイントからログの消費を開始します。コンシューマーグループに対応するチェックポイントがない場合、消費は startTime の設定値から開始されます。

    • false (デフォルト値):指定されたコンシューマーグループに保存されているチェックポイントからログの消費を開始しません。

    重要

    このパラメーターは VVR 11.1 以降ではサポートされなくなりました。これらのバージョンでは、startupMode オプションを consumer_group に設定する必要があります。

    maxRetries

    SLS からの読み取り試行が失敗した後の再試行回数。

    String

    いいえ

    3

    N/A

    batchGetSize

    1 回のリクエストで読み取るロググループの数。

    String

    いいえ

    100

    batchGetSize の設定は 1000 を超えることはできません。超えるとエラーが報告されます。

    exitAfterFinish

    すべてのデータが消費された後に Flink ジョブが終了するかどうかを指定します。

    String

    いいえ

    false

    • true:すべてのデータが消費された後、Flink プログラムは終了します。

    • false (デフォルト):データ消費が完了しても Flink プログラムは終了しません。

    query

    重要

    このオプションは VVR 11.3 で非推奨になりましたが、以降のバージョンでも互換性は維持されています。

    消費前にデータを前処理するためのクエリ文。

    String

    いいえ

    なし

    このオプションを使用して、Flink が消費する前に SLS データをフィルタリングします。これにより、コストが削減され、処理速度が向上します。

    例えば、 'query' = '*| where request_method = ''GET''' は、Flink が SLS からデータを読み取る前に、まず request_method フィールドの値が 'GET' であるデータを照合することを示します。

    説明

    このオプションは Log Service (SLS) の SPL 言語を使用します。詳細については、「SPL 構文」をご参照ください。

    重要
    • このオプションは VVR 8.0.1 以降でのみサポートされています。

    • この機能には Log Service (SLS) の料金が発生します。詳細については、「Log Service の課金」をご参照ください。

    processor

    データの前処理に使用する SLS プロセッサの名前。queryprocessor の両方が指定されている場合、query が優先され、processor は無視されます。

    String

    いいえ

    なし

    このオプションは、Flink が消費する前に SLS データをフィルタリングするため、コストを削減し、処理速度を向上させます。query の代わりに processor を使用することを推奨します。

    例えば、 'processor' = 'test-filter-processor' は、Flink が SLS からデータを読み取る前に、SLS プロセッサがデータをフィルタリングすることを示します。

    説明

    このオプションは Log Service (SLS) の SPL 言語を使用します。詳細については、「SPL 構文」をご参照ください。SLS プロセッサの作成または更新方法については、「プロセッサの管理」をご参照ください。

    重要

    このオプションは VVR 11.3 以降でのみサポートされています。

    この機能には Log Service (SLS) の料金が発生します。詳細については、「Log Service の課金」をご参照ください。

    preserveRawBytes

    enableNewSourcetrue に設定されている場合に、SLS が運ぶ生のバイトを直接フェッチして保持するかどうかを指定します。

    Boolean

    いいえ

    false

    このオプションを有効にすると、BINARY/VARBINARY フィールドは生の byte[] を直接読み取り、CHAR/VARCHAR フィールドは生のバイトから文字列データを構築します。他のフィールドは、引き続きデフォルトの変換ロジックを使用します。この動作は、古いソースの動作と一致します。

    説明
    • このオプションは VVR 11.8 以降でのみサポートされています。

    • このオプションは、enableNewSourcetrue に設定されている場合にのみ有効です。

  • sink 固有

    パラメーター

    説明

    タイプ

    必須

    デフォルト

    備考

    topicField

    ログトピックを示す __topic__ 属性を上書きする値を指定するフィールド。

    String

    いいえ

    なし

    このオプションの値は、テーブル内の既存のフィールドでなければなりません。

    timeField

    ログ書き込み時間を示す __timestamp__ 属性を上書きする値を指定するフィールド。

    String

    いいえ

    現在時刻

    このオプションの値は、テーブル内の既存の INT フィールドでなければなりません。このオプションが指定されていない場合、現在時刻が使用されます。

    sourceField

    ログソース (ログを生成したマシンの IP アドレスなど) を示す __source__ 属性を上書きする値を指定するフィールド。

    String

    いいえ

    なし

    このオプションの値は、テーブル内の既存のフィールドでなければなりません。

    partitionField

    パーティショニング用のフィールドを指定します。このフィールドの値のハッシュによって、どのシャードがデータを受信するかを決定し、同じハッシュを持つレコードが同じシャードに送られるようにします。

    String

    いいえ

    なし

    このオプションが指定されていない場合、各レコードは利用可能なシャードにランダムに書き込まれます。

    buckets

    partitionField が指定されている場合、このオプションはハッシュ値をマッピングするためのバケット数を定義します。

    String

    いいえ

    64

    値は [1, 256] の範囲の 2 のべき乗でなければなりません。バケット数はシャード数以上でなければなりません。そうでない場合、一部のシャードはデータを受信しない可能性があります。

    flushIntervalMs

    データ書き込み間隔。

    String

    いいえ

    2000

    単位:ミリ秒。

    writeNullProperties

    null 値を空の文字列として SLS に書き込むかどうかを指定します。

    Boolean

    いいえ

    true

    • true (デフォルト値):null 値をログに空の文字列として書き込みます。

    • false:null と評価されるフィールドはログに書き込まれません。

    説明

    このオプションは VVR 8.0.6 以降でのみサポートされています。

型マッピング

Flink 型

SLS 型

BOOLEAN

STRING

VARBINARY

VARCHAR

TINYINT

INTEGER

BIGINT

FLOAT

DOUBLE

DECIMAL

データインジェスト (ベータ版)

制限事項

この機能は、Realtime Compute for Apache Flink バージョン 11.1 以降でのみサポートされています。

構文

source:
   type: sls
   name: SLS Source
   endpoint: <endpoint>
   project: <project>
   logstore: <logstore>
   accessId: <accessId>
   accessKey: <accessKey>

パラメーター

パラメーター

説明

タイプ

必須

デフォルト

備考

type

データソースのタイプ。

String

はい

なし

値は sls でなければなりません。

endpoint

Log Service (SLS) のエンドポイント。

String

はい

なし

Log Service (SLS) の VPC アクセスアドレス。詳細については、「サービスエンドポイント」をご参照ください。

説明
  • デフォルトでは、Realtime Compute for Apache Flink はインターネットにアクセスできません。NAT Gateway を使用して、ご利用の Virtual Private Cloud (VPC) とインターネット間の通信を有効にできます。詳細については、「インターネットへのアクセス方法」をご参照ください。

  • インターネット経由で Log Service (SLS) にアクセスすることは推奨されません。もしそうする必要がある場合は、HTTPS を使用し、SLS の 転送アクセラレーション を有効にしてください。

accessId

ご利用の Alibaba Cloud アカウントの AccessKey ID。

String

はい

なし

詳細については、「AccessKey ID と AccessKey Secret 情報を表示する方法」をご参照ください。

重要

AccessKey 情報が公開されるのを防ぐため、プロジェクト変数を使用して AccessKey の値を指定することを推奨します。詳細については、「プロジェクト変数」をご参照ください。

accessKey

ご利用の Alibaba Cloud アカウントの AccessKey Secret。

String

はい

なし

project

Log Service (SLS) プロジェクトの名前。

String

はい

なし

なし

logStore

SLS Logstore または Metricstore の名前。

String

はい

なし

Logstore 内のデータは、Metricstore 内のデータと同じ方法で消費されます。

schema.inference.strategy

スキーマ推論戦略。

String

いいえ

continuous

  • continuous:各データレコードに対してスキーマ推論を実行します。スキーマに互換性がない場合、より広いスキーマが推論され、スキーマ変更イベントが生成されます。

  • static:ジョブ開始時に一度だけスキーマ推論を実行します。後続のデータは初期スキーマに基づいて解析され、スキーマ変更イベントは生成されません。

maxPreFetchLogGroups

初期スキーマ推論のために各シャードから読み取るロググループの最大数。

Integer

いいえ

50

ジョブがデータを読み取って処理する前に、コネクタは各シャードから指定された数のロググループを事前に消費して、スキーマ情報を初期化します。

shardDiscoveryIntervalMs

シャードの変更を動的に検出する間隔 (ミリ秒単位)。

Long

いいえ

60000

このパラメーターに負の値を設定すると、動的検出が無効になります。

説明

値は 60,000 ミリ秒 (1 分) 以上である必要があります。

startupMode

起動モード。

String

いいえ

なし

  • timestamp (デフォルト):特定のタイムスタンプからログを消費します。

  • latest:最新のオフセットからログを消費します。

  • earliest:最も古いオフセットからログを消費します。

  • consumer_group:コンシューマーグループに記録されたオフセットからログを消費します。シャードにオフセットが記録されていない場合、消費は最も古いオフセットから開始されます。

startTime

ログ消費の開始時刻。

String

いいえ

現在時刻

フォーマットは yyyy-MM-dd HH:mm:ss です。

このパラメーターは、startupModetimestamp に設定されている場合にのみ有効です。

説明

startTime および stopTime パラメーターは、Log Service (SLS) の __timestamp__ 属性ではなく、__receive_time__ 属性に基づいています。

stopTime

ログ消費の終了時刻。

String

いいえ

なし

フォーマットは yyyy-MM-dd HH:mm:ss です。

説明

すべてのログが消費された後に Flink ジョブを終了させたい場合は、exitAfterFinish=true も設定する必要があります。

consumerGroup

コンシューマーグループの名前。

String

いいえ

なし

コンシューマーグループは消費の進行状況を記録します。任意のカスタム名を指定できます。

batchGetSize

リクエストごとに読み取るロググループの数。

Integer

いいえ

100

batchGetSize の値は 1,000 を超えることはできません。超えるとエラーが発生します。

maxRetries

Log Service (SLS) からの読み取りが失敗した場合の再試行回数。

Integer

いいえ

3

なし

exitAfterFinish

すべてのデータが消費された後に Flink ジョブが終了するかどうかを指定します。

Boolean

いいえ

false

  • true:すべてのデータが消費された後、Flink ジョブは終了します。

  • false (デフォルト):すべてのデータが消費された後も Flink ジョブは終了しません。

query

Log Service (SLS) からデータを消費するための前処理文。

String

いいえ

なし

このパラメーターを使用して、消費前に Log Service (SLS) のデータをフィルタリングし、コストを節約し、処理速度を向上させます。

例えば、'query' = '*| where request_method = ''GET''' は、Flink によってデータが読み取られる前に、request_method フィールドが 'GET' であるデータをフィルタリングします。

説明

クエリは Log Service の SPL 構文を使用する必要があります。詳細については、「SPL 構文」をご参照ください。

重要
  • この機能が利用可能な Log Service (SLS) のリージョンについては、「ルールに基づいてログを消費する」をご参照ください。

  • この機能はベータ版であり、無料です。将来的にこの機能に対して課金される可能性があります。詳細については、「料金」をご参照ください。

compressType

Log Service (SLS) の圧縮タイプ。

String

いいえ

なし

サポートされている圧縮タイプは次のとおりです:

  • lz4

  • deflate

  • zstd

timeZone

startTimestopTime のタイムゾーン。

String

いいえ

なし

デフォルトでは、オフセットは追加されません。

regionId

Log Service (SLS) がデプロイされているリージョン。

String

いいえ

なし

詳細については、「サポートされているリージョン」をご参照ください。

signVersion

Log Service (SLS) のリクエスト署名バージョン。

String

いいえ

なし

詳細については、「リクエスト署名」をご参照ください。

shardModDivisor

Logstore シャードから読み取る際に使用される除数。

Int

いいえ

-1

詳細については、「シャード」をご参照ください。

shardModRemainder

Logstore シャードから読み取る際に使用される剰余。

Int

いいえ

-1

詳細については、「シャード」をご参照ください。

metadata.list

下流のジョブに渡すメタデータ列。

String

いいえ

なし

利用可能なメタデータフィールドには、__source____topic____timestamp__、および __tag__ があります。複数のフィールドはカンマで区切ります。

decode.table-id.fields

Log Service (SLS) からのログデータを解析する際に、テーブル ID を生成するために使用される値を持つフィールドを指定します。

String

いいえ

なし

複数のフィールドは英語のカンマ , で区切られます。例えば、上流の SLS ログレコードが {"col0":"a", "col1":"b", "col2":"c"} の場合、異なるパラメーター構成の結果は次のようになります:

構成

テーブル ID

なし

すべてのメッセージは Project.Logstore です

col0

a

col0,col1

a.b

col0,col1,col2

a.b.c

説明

このパラメーターは、Realtime Compute for Apache Flink バージョン 11.6 以降でサポートされています。

fixed-types

Log Service (SLS) からのログデータを解析する際に、特定のフィールドのデータ型を指定します。

String

いいえ

なし

データを解析する際に、特定のフィールドの型を指定します。複数のフィールド定義を区切るには、カンマ , を使用します。例えば、id BIGINT, name VARCHAR(10) は、id フィールドの型を BIGINT に、name フィールドの型を VARCHAR(10) に指定します。

説明

このパラメーターは、Realtime Compute for Apache Flink バージョン 11.6 以降でサポートされています。

timestamp-format.standard

Log Service (SLS) からのログデータ内のタイムスタンプフィールドのフォーマット。

String

いいえ

SQL

有効な値:

  • SQL:入力タイムスタンプを yyyy-MM-dd HH:mm:ss.s{precision} フォーマット (例:2020-12-30 12:13:14.123) で解析し、同じフォーマットで出力します。

  • ISO-8601:入力タイムスタンプを yyyy-MM-ddTHH:mm:ss.s{precision} フォーマット (例:2020-12-30T12:13:14.123) で解析し、同じフォーマットで出力します。

説明

このパラメーターは、Realtime Compute for Apache Flink バージョン 11.6 以降でサポートされています。

ingestion.ignore-errors

データ解析中に発生したエラーを無視するかどうかを指定します。

Boolean

いいえ

false

説明

このパラメーターは、Realtime Compute for Apache Flink バージョン 11.6 以降でサポートされています。

ingestion.error-tolerance.max-count

ingestion.ignore-errors が有効な場合、累積エラー数がこの値を超えるとジョブは失敗します。

Integer

いいえ

-1

このパラメーターは、ingestion.ignore-errors が有効な場合にのみ有効です。デフォルト値の -1 は、ジョブがすべての解析例外を無視することを意味します。

説明

このパラメーターは、Realtime Compute for Apache Flink バージョン 11.6 以降でサポートされています。

既存のカタログの再利用

Realtime Compute for Apache Flink バージョン 11.5 以降、Flink CDC データインジェストジョブで Data Management ページに作成された組み込み SLS カタログを参照できます。これにより、接続プロパティを手動で指定する必要がなくなります。

source:
  type: sls
  using.built-in-catalog: sls_catalog

現在、データインジェストジョブは、組み込み SLS カタログから次のパラメーターを自動的に再利用できます:

  • endpoint

  • project

  • accessId

  • accessKey

これらの自動的に再利用されるパラメーターを上書きするには、YAML 構成で明示的に定義できます。YAML ファイルで定義されたパラメーターが優先されます。

データ型マッピング

fixed-types が構成されていない場合、次のデータ型マッピングが適用されます:

SLS 型

CDC 型

STRING

STRING

fixed-types が構成されている場合、システムは指定された型を使用してデータを解析します。

スキーマの推論と進化

  • シャードデータの事前消費とスキーマの初期化

    SLS コネクタは、読み取っている Logstore のスキーマを維持します。Logstore からデータを読み取る前に、コネクタは各シャードから最大 maxPreFetchLogGroups 個のロググループを事前に消費します。各ログエントリのスキーマを解析し、それらをマージしてテーブルのスキーマを初期化します。データ消費が始まる前に、この初期スキーマに基づいてテーブル作成イベントが生成されます。

    説明

    各シャードについて、コネクタは現在時刻の 1 時間前からデータの消費を開始してログスキーマを解析しようとします。

  • プライマリキー情報

    Log Service (SLS) のログにはプライマリキー情報が含まれていません。変換ルールを使用して、テーブルに手動でプライマリキーを追加できます:

    transform:
      - source-table: <project>.<logstore>
        projection: *
        primary-keys: key1, key2
  • スキーマの推論とスキーマの変更

    スキーマが初期化された後、schema.inference.strategystatic に設定されている場合、SLS コネクタは初期スキーマに基づいて各ログエントリを解析し、スキーマ変更イベントを生成しません。schema.inference.strategycontinuous に設定されている場合、コネクタは各ログエントリを解析し、物理列を推論し、それらを現在のスキーマと比較します。推論されたスキーマが現在のスキーマと一致しない場合、スキーマは次のルールに従ってマージされます:

    • 推論された物理列に現在のスキーマにないフィールドが含まれている場合、コネクタはこれらのフィールドをスキーマに追加し、null 許容列を追加するイベントを生成します。

    • 推論された物理列に現在のスキーマに存在するフィールドが欠けている場合、コネクタはこれらのフィールドを保持し、そのデータを NULL で埋め、列削除イベントを生成しません。

    SLS コネクタは、各ログエントリのすべてのフィールドのデータ型を String として推論します。現在、新しい列の追加のみがサポートされています。コネクタは、新しい列をスキーマの末尾に追加し、null 許容に設定します。

コード例

  • ソーステーブルと結果テーブルの SQL

    CREATE TEMPORARY TABLE sls_input(
      `time` BIGINT,
      url STRING,
      dt STRING,
      float_field FLOAT,
      double_field DOUBLE,
      boolean_field BOOLEAN,
      `__topic__` STRING METADATA VIRTUAL,
      `__source__` STRING METADATA VIRTUAL,
      `__timestamp__` STRING METADATA VIRTUAL,
       __tag__ MAP<VARCHAR, VARCHAR> METADATA VIRTUAL,
      proctime as PROCTIME()
    ) WITH (
      'connector' = 'sls',
      'endpoint' ='cn-hangzhou-intranet.log.aliyuncs.com',
      'accessId' = '${secret_values.ak_id}',
      'accessKey' = '${secret_values.ak_secret}',
      'starttime' = '2023-08-30 00:00:00',
      'project' ='sls-test',
      'logstore' ='sls-input'
    );
    
    CREATE TEMPORARY TABLE sls_sink(
      `time` BIGINT,
      url STRING,
      dt STRING,
      float_field FLOAT,
      double_field DOUBLE,
      boolean_field BOOLEAN,
      `__topic__` STRING,
      `__source__` STRING,
      `__timestamp__` BIGINT ,
      receive_time BIGINT
    ) WITH (
      'connector' = 'sls',
      'endpoint' ='cn-hangzhou-intranet.log.aliyuncs.com',
      'accessId' = '${ak_id}',
      'accessKey' = '${ak_secret}',
      'project' ='sls-test',
      'logstore' ='sls-output'
    );
    
    INSERT INTO sls_sink
    SELECT 
     `time`,
      url,
      dt,
      float_field,
      double_field,
      boolean_field,
      `__topic__` ,
      `__source__` ,
      `__timestamp__` ,
      cast(__tag__['__receive_time__'] as bigint) as receive_time
    FROM sls_input; 
  • SLS データソースを使用したデータインジェスト

    SLS をデータソースとして使用して、サポートされている下流システムにリアルタイムでデータをインジェストします。例えば、次の構成では、Logstore から Data Lake Formation (DLF) の Paimon 形式のデータレイクにデータを書き込むデータインジェストジョブを定義します。ジョブは結果テーブルのスキーマを自動的に推論し、ランタイムでのスキーマ進化をサポートします。

source:
  type: sls
  name: SLS Source
  endpoint: ${endpoint}
  project: test_project
  logstore: test_log
  accessId: ${accessId}
  accessKey: ${accessKey}
   
# テーブルにプライマリキーを追加します。
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id
    
# test_project.test_log からのすべてのデータを test_database.inventory テーブルにルーティングします。
route:
  - source-table: test_project.test_log
    sink-table: test_database.inventory

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (オプション) 読み取りパフォーマンスを向上させるために削除ベクターを有効にします。
  table.properties.deletion-vectors.enabled: true

DataStream API

重要

DataStream API でデータを読み書きするには、DataStream コネクタを使用します。詳細については、「DataStream コネクタの使用法」をご参照ください。

VVR バージョン 8.0.10 より前を使用している場合、依存関係の欠落によりジョブの開始に失敗することがあります。この問題を解決するには、対応する uber-JAR を追加の依存関係として追加してください。

SLS からの読み取り

Realtime Compute for Apache Flink は、Simple Log Service (SLS) からデータを読み取るための SourceFunction の実装である SlsSourceFunction クラスを提供します。次の例では、SLS からデータを読み取ります。

public class SlsDataStreamSource {

    public static void main(String[] args) throws Exception {
        // ストリーミング実行環境を設定します
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        // SLS ソースを作成し、データをコンソールに出力します。
        env.addSource(createSlsSource())
                .map(SlsDataStreamSource::convertMessages)
                .print();
        env.execute("SLS Stream Source");
    }

    private static SlsSourceFunction createSlsSource() {
        SLSAccessInfo accessInfo = new SLSAccessInfo();
        accessInfo.setEndpoint("yourEndpoint");
        accessInfo.setProjectName("yourProject");
        accessInfo.setLogstore("yourLogStore");
        accessInfo.setAccessId("yourAccessId");
        accessInfo.setAccessKey("yourAccessKey");

        // バッチ取得サイズは必須です。
        accessInfo.setBatchGetSize(10);

        // オプションのパラメーター
        accessInfo.setConsumerGroup("yourConsumerGroup");
        accessInfo.setMaxRetries(3);

        // 消費の開始時刻、現在時刻に設定。
        int startInSec = (int) (new Date().getTime() / 1000);

        // 消費の停止時刻、-1 は停止しないことを意味します。
        int stopInSec = -1;

        return new SlsSourceFunction(accessInfo, startInSec, stopInSec);
    }

    private static List<String> convertMessages(SourceRecord input) {
        List<String> res = new ArrayList<>();
        for (FastLogGroup logGroup : input.getLogGroups()) {
            int logsCount = logGroup.getLogsCount();
            for (int i = 0; i < logsCount; i++) {
                FastLog log = logGroup.getLogs(i);
                int fieldCount = log.getContentsCount();
                for (int idx = 0; idx < fieldCount; idx++) {
                    FastLogContent f = log.getContents(idx);
                    res.add(String.format("key: %s, value: %s", f.getKey(), f.getValue()));
                }
            }
        }
        return res;
    }
}

SLS への書き込み

Realtime Compute for Apache Flink は、SLS にデータを書き込むための OutputFormat の実装である SLSOutputFormat クラスを提供します。次の例では、SLS にデータを書き込みます。

public class SlsDataStreamSink {

    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.fromSequence(0, 100)
                .map((MapFunction<Long, SinkRecord>) aLong -> getSinkRecord(aLong))
                .addSink(createSlsSink())
                .name(SlsDataStreamSink.class.getSimpleName());
        env.execute("SLS Stream Sink");
    }

    private static OutputFormatSinkFunction createSlsSink() {
        Configuration conf = new Configuration();
        conf.setString(SLSOptions.ENDPOINT, "yourEndpoint");
        conf.setString(SLSOptions.PROJECT, "yourProject");
        conf.setString(SLSOptions.LOGSTORE, "yourLogStore");
        conf.setString(SLSOptions.ACCESS_ID, "yourAccessId");
        conf.setString(SLSOptions.ACCESS_KEY, "yourAccessKey");
        SLSOutputFormat outputFormat = new SLSOutputFormat(conf);
        return new OutputFormatSinkFunction<>(outputFormat);
    }

    private static SinkRecord getSinkRecord(Long seed) {
        SinkRecord record = new SinkRecord();
        LogItem logItem = new LogItem((int) (System.currentTimeMillis() / 1000));
        logItem.PushBack("level", "info");
        logItem.PushBack("name", String.valueOf(seed));
        logItem.PushBack("message", "it's a test message for " + seed.toString());
        record.setContent(logItem);
        return record;
    }

}

XML

SLS DataStream コネクタは、Maven セントラルリポジトリで利用可能です。

<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-sls</artifactId>
    <version>${vvr-version}</version>
    <exclusions>
        <exclusion>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-format-common</artifactId>
        </exclusion>
    </exclusions>
</dependency>

よくある質問

失敗した Flink プログラムを復元する際の TaskManager の OOM (java.lang.OutOfMemoryError: Java heap space) を解決する方法は?