Kafka 이벤트 소스용 오류 처리 제어 구성하기
Kafka 이벤트 소스용 오류 처리 제어 구성하기 (Configuring error handling controls for Kafka event sources)
Kafka 이벤트 소스 매핑에서 Lambda가 오류와 재시도를 처리하는 방법을 구성할 수 있어요. 이 구성은 Lambda가 실패한 레코드를 처리하고 재시도 동작을 관리하는 방식을 제어하도록 도와줘요.
사용 가능한 재시도 구성
Amazon MSK와 자체 관리 Kafka 이벤트 소스 모두에 다음 재시도 구성이 제공돼요.
- 최대 재시도 횟수(Maximum retry attempts) — 함수가 오류를 반환할 때 Lambda가 재시도하는 최대 횟수예요. 초기 호출 시도는 포함하지 않아요. 기본값은 -1(무제한)이에요. 무제한 재시도와 실패 시 대상(on-failure destination)을 함께 구성하면 Lambda가 자동으로 최대 10회 재시도를 적용해요.
- 최대 레코드 보존 기간(Maximum record age) — Lambda가 함수에 보내는 레코드의 최대 보존 기간이에요. 기본값은 -1(무제한)이에요.
- 오류 시 배치 분할(Split batch on error) — 함수가 오류를 반환하면 배치를 두 개의 더 작은 배치로 나누고 각각을 별도로 재시도해요. 이 기능은 문제가 있는 레코드를 격리하는 데 도움을 줘요.
- 부분 배치 응답(Partial batch response) — 함수가 배치에서 처리가 실패한 레코드에 대한 정보를 반환하게 해서, Lambda가 실패한 레코드만 재시도할 수 있게 해줘요.
오류 처리 제어 구성하기 (콘솔)
Lambda 콘솔에서 Kafka 이벤트 소스 매핑을 만들거나 업데이트할 때 재시도 동작을 구성할 수 있어요.
Kafka 이벤트 소스의 재시도 동작 구성하기 (콘솔)
- Lambda 콘솔의 Functions 페이지를 열어요.
- 함수 이름을 선택해요.
- 다음 중 하나를 해요. 새 Kafka 트리거를 추가하려면 Function overview 아래에서 Add trigger를 선택해요. 기존 Kafka 트리거를 수정하려면 트리거를 선택한 다음 Edit을 선택해요.
- Event poller configuration 아래에서 프로비저닝 모드를 선택해 오류 처리 제어를 구성해요.
- Retry attempts에 최대 재시도 횟수(0~10000, 또는 무제한은 -1)를 입력해요.
- Maximum record age에 최대 보존 기간을 초 단위(60~604800, 또는 무제한은 -1)로 입력해요.
- 오류 발생 시 배치 분할을 활성화하려면 Split batch on error를 선택해요.
- 부분 배치 응답을 활성화하려면 ReportBatchItemFailures를 선택해요.
- Add 또는 Save를 선택해요.
재시도 동작 구성하기 (AWS CLI)
Kafka 이벤트 소스 매핑의 재시도 동작을 구성하려면 다음 AWS CLI 명령을 사용해요.
재시도 구성과 함께 이벤트 소스 매핑 만들기
다음 예시는 오류 처리 제어와 함께 자체 관리 Kafka 이벤트 소스 매핑을 만들어요.
aws lambda create-event-source-mapping \
--function-name my-kafka-function \
--topics my-kafka-topic \
--source-access-configuration Type=SASL_SCRAM_512_AUTH,URI=arn:aws:secretsmanager:us-east-1:111122223333:secret:MyBrokerSecretName \
--self-managed-event-source '{"Endpoints":{"KAFKA_BOOTSTRAP_SERVERS":["abc.xyz.com:9092"]}}' \
--starting-position LATEST \
--provisioned-poller-config MinimumPollers=1,MaximumPollers=1 \
--maximum-retry-attempts 3 \
--maximum-record-age-in-seconds 3600 \
--bisect-batch-on-function-error \
--function-response-types "ReportBatchItemFailures"
Amazon MSK 이벤트 소스의 경우:
aws lambda create-event-source-mapping \
--event-source-arn arn:aws:kafka:us-east-1:111122223333:cluster/my-cluster/fc2f5bdf-fd1b-45ad-85dd-15b4a5a6247e-2 \
--topics AWSMSKKafkaTopic \
--starting-position LATEST \
--function-name my-kafka-function \
--source-access-configurations '[{"Type": "SASL_SCRAM_512_AUTH","URI": "arn:aws:secretsmanager:us-east-1:111122223333:secret:my-secret"}]' \
--provisioned-poller-config MinimumPollers=1,MaximumPollers=1 \
--maximum-retry-attempts 3 \
--maximum-record-age-in-seconds 3600 \
--bisect-batch-on-function-error \
--function-response-types "ReportBatchItemFailures"
재시도 구성 업데이트하기
update-event-source-mapping 명령으로 기존 이벤트 소스 매핑의 재시도 구성을 수정해요.
aws lambda update-event-source-mapping \
--uuid 12345678-1234-1234-1234-123456789012 \
--maximum-retry-attempts 5 \
--maximum-record-age-in-seconds 7200 \
--bisect-batch-on-function-error \
--function-response-types "ReportBatchItemFailures"
PartialBatchResponse
부분 배치 응답, 다른 말로 ReportBatchItemFailures는 Lambda와 Kafka 소스 통합에서 오류 처리의 핵심 기능이에요. 이 기능이 없으면 배치의 항목 하나에서 오류가 발생했을 때 그 배치의 모든 메시지가 다시 처리돼요. 부분 배치 응답을 활성화하고 구현하면 핸들러가 실패한 메시지의 식별자만 반환해서, Lambda가 그 특정 항목들만 재시도할 수 있게 해줘요. 이렇게 하면 실패한 메시지가 포함된 배치가 처리되는 방식을 더 잘 제어할 수 있어요.
배치 오류를 보고하려면 이 JSON 스키마를 사용해요.
{
"batchItemFailures": [
{
"itemIdentifier": {
"partition": "topic-partition_number",
"offset": 100
}
},
...
]
}
중요 빈 유효 JSON이나 null을 반환하면 이벤트 소스 매핑은 배치가 성공적으로 처리된 것으로 간주해요. 호출된 이벤트에 없던 잘못된 topic-partition_number 또는 offset을 반환하면 실패로 처리되고 전체 배치가 재시도돼요.
다음 코드 예시는 Kafka 소스에서 이벤트를 받는 Lambda 함수에 부분 배치 응답을 구현하는 방법을 보여줘요. 함수는 응답에서 배치 항목 실패를 보고해서, Lambda가 나중에 그 메시지들을 재시도하도록 신호를 보내요.
다음은 이 방식을 보여주는 Python Lambda 핸들러 구현이에요.
import base64
from typing import Any, Dict, List
def lambda_handler(event: Dict[str, Any], context: Any) -> Dict[str, List[Dict[str, Dict[str, Any]]]]:
failures: List[Dict[str, Dict[str, Any]]] = []
records_dict = event.get("records", {})
for topic_partition, records_list in records_dict.items():
for record in records_list:
topic = record.get("topic")
partition = record.get("partition")
offset = record.get("offset")
value_b64 = record.get("value")
try:
data = base64.b64decode(value_b64).decode("utf-8")
process_message(data)
except Exception as exc:
print(f"Failed to process record topic={topic} partition={partition} offset={offset}: {exc}")
item_identifier: Dict[str, Any] = {
"partition": f"{topic}-{partition}",
"offset": int(offset) if offset is not None else None,
}
failures.append({"itemIdentifier": item_identifier})
return {"batchItemFailures": failures}
def process_message(data: str) -> None:
# Your business logic for a single message
pass
Node.js 버전:
const { Buffer } = require("buffer");
const handler = async (event) => {
const failures = [];
for (let topicPartition in event.records) {
const records = event.records[topicPartition];
for (const record of records) {
const topic = record.topic;
const partition = record.partition;
const offset = record.offset;
const valueBase64 = record.value;
const data = Buffer.from(valueBase64, "base64").toString("utf8");
try {
await processMessage(data);
} catch (error) {
console.error("Failed to process record", { topic, partition, offset, error });
const itemIdentifier = {
"partition": `${topic}-${partition}`,
"offset": Number(offset),
};
failures.push({ itemIdentifier });
}
}
}
return { batchItemFailures: failures };
};
async function processMessage(payload) {
// Your business logic for a single message
}
module.exports = { handler };
Java 버전:
import com.amazonaws.services.lambda.runtime.Context;
import com.amazonaws.services.lambda.runtime.RequestHandler;
import java.util.ArrayList;
import java.util.Base64;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
public class KafkaBatchHandler implements RequestHandler<Map<String, Object>, Map<String, Object>> {
@SuppressWarnings("unchecked")
@Override
public Map<String, Object> handleRequest(Map<String, Object> event, Context context) {
List<Map<String, Object>> failures = new ArrayList<>();
Map<String, List<Map<String, Object>>> records =
(Map<String, List<Map<String, Object>>>) event.getOrDefault("records", Map.of());
for (Map.Entry<String, List<Map<String, Object>>> entry : records.entrySet()) {
for (Map<String, Object> record : entry.getValue()) {
String topic = (String) record.get("topic");
Object partition = record.get("partition");
Object offset = record.get("offset");
String valueBase64 = (String) record.get("value");
try {
String data = new String(Base64.getDecoder().decode(valueBase64), "UTF-8");
processMessage(data);
} catch (Exception e) {
System.err.printf("Failed to process record topic=%s partition=%s offset=%s: %s%n",
topic, partition, offset, e.getMessage());
Map<String, Object> itemIdentifier = new HashMap<>();
itemIdentifier.put("partition", topic + "-" + partition);
itemIdentifier.put("offset", offset instanceof Number ? ((Number) offset).longValue() : null);
Map<String, Object> failure = new HashMap<>();
failure.put("itemIdentifier", itemIdentifier);
failures.add(failure);
}
}
}
Map<String, Object> response = new HashMap<>();
response.put("batchItemFailures", failures);
return response;
}
private void processMessage(String data) {
// Your business logic for a single message
}
}