You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

从Amazon Connect经Kinesis Video Stream提取的直播音频为何嘈杂残缺?

音频嘈杂卡顿且内容缺失问题排查

问题背景

基于Amazon Connect官方示例编写Lambda函数,将Kinesis Video Stream(KVS)中的直播音频块推送至外部WebSocket服务器。音频规格为16-bit PCM单声道、8kHz采样率。

接收端处理流程:

  • 将每个WebSocket收到的音频块保存为带唯一时间戳的.raw文件
  • 使用sox将原始音频转换为WAV文件:
$ sox -r 8k -c 2 -e signed -b16 -t raw chunk-1692154346612.raw chunk-1692154346612.wav
  • 合并所有WAV文件:
$ sox *.wav all-chunks.wav

Shell扩展*.wav会按时间顺序排序文件,但将合并后的all-chunks.wav导入Audacity播放时,音频整体嘈杂卡顿,且大量通话内容缺失。

Lambda处理器逻辑

public class ProcessAudioOverKVS implements RequestHandler<TranscriptionRequest, String>
{
  private static final Logger logger = LoggerFactory.getLogger(ProcessAudioOverKVS.class);
  private static final Regions REGION = Regions.fromName(System.getenv("APP_REGION"));
  private static final String START_SELECTOR_TYPE = System.getenv("START_SELECTOR_TYPE");
  private static final int CHUNK_SIZE_IN_KB = 4;
  private static final String RELAY_AUDIO_CHUNK_URL = System.getenv("RELAY_AUDIO_CHUNK_URL");

  /**
   * Handler function for the Lambda. TranscriptionRequest has these properties:
   * - streamARN: unique ARN that identifies live audio stream over Kinesis Video Streams.
   * - startFragmentNum: integer denoting start of audio content.
   * - connectContactId: unique ID that identifies the caller.
   * - transcriptionEnabled: boolean, always false.
   * - languageCode: set to en-US.
   * - saveCallRecording: boolean, always set to false.
   * - streamAudioFromCustomer: boolean, always true.
   * - streamAudioToCustomer: boolean, always true,
   *
   * @param request
   * @param context
   * @return
   */
  @Override
  public String handleRequest(TranscriptionRequest request, Context context)
  {
    logger.info("received request : " + request.toString());
    logger.info("received context: " + context.toString());

    boolean success = false;
    WebSocketClient wsClient = new WebSocketClient();
    try {
      request.validate();
      wsClient.Connect(RELAY_AUDIO_CHUNK_URL);
      pushAudioChunks(wsClient, request);
      success = true;
    } catch (Exception e) {
      logger.error("ProcessAudioOverKVS failed", e);
    } finally {
      try { wsClient.Disconnect(); } catch (Exception e) {}
    }
    return success ? "{ \"result\": \"Success\" }" : "{ \"result\": \"Failed\" }";
  }

  private void pushAudioChunks(WebSocketClient wsClient, TranscriptionRequest request) throws Exception
  {
    logger.debug("pushAudioChunks ...");
    String streamARN = request.getStreamARN();
    String streamName = streamARN.substring(streamARN.indexOf("/") + 1, streamARN.lastIndexOf("/"));
    String contactId = request.getConnectContactId();
    String startFragmentNum = request.getStartFragmentNum();

    InputStream kvsInputStream = KVSUtils.getInputStreamFromKVS(
      streamName, REGION, startFragmentNum, getAWSCredentials(), START_SELECTOR_TYPE
    );
    StreamingMkvReader streamingMkvReader = StreamingMkvReader.createDefault(new InputStreamParserByteSource(kvsInputStream));
    KVSContactTagProcessor tagProcessor = new KVSContactTagProcessor(contactId);
    FragmentMetadataVisitor fragmentVisitor = FragmentMetadataVisitor.create(Optional.of(tagProcessor));

    List<KVSUtils.AudioChunkTrack> chunkList = KVSUtils.getByteBufferFromStream(
      streamingMkvReader, fragmentVisitor, tagProcessor, contactId, CHUNK_SIZE_IN_KB
    );
    boolean haveAudio = chunkList.size() > 0 && chunkList.get(0).audio.remaining() > 0;
    while (haveAudio) {
      for (KVSUtils.AudioChunkTrack act : chunkList) {
        if (act.audio.remaining() > 0) {
          byte[] audioBytes = new byte[act.audio.remaining()];
          act.audio.get(audioBytes);
          try {
            relayAudioChunk(wsClient, contactId, audioBytes, act.track, act.trackNum);
          } catch (Exception ex) {
            logger.error(
              "could not relay audio chunk for contactId={}, track={}, trackNum={}",
              contactId, act.track, act.trackNum, ex
            );
          }
        }
      }
      chunkList = KVSUtils.getByteBufferFromStream(
        streamingMkvReader, fragmentVisitor, tagProcessor, contactId, CHUNK_SIZE_IN_KB
      );
      haveAudio = chunkList.size() > 0 && chunkList.get(0).audio.remaining() > 0;
    }
  }

  /**
   * Relay audio chunk to Minerva.
   */
  private void relayAudioChunk(
    WebSocketClient wsClient,
    String contactId,
    byte[] audioBytes,
    String track,
    long trackNum
  ) throws Exception {
    logger.debug("relayAudioChunk ...");
    AudioChunk audioChunk = new AudioChunk();
    audioChunk.setContactId(contactId);
    audioChunk.setTrackName(track);
    audioChunk.setTrackNum(trackNum);
    audioChunk.setAudio(audioBytes);
    wsClient.SendAudioChunk(audioChunk);
  }

  /**
   * @return AWS credentials to be used to connect to Transcribe service. This example uses the default credentials
   * provider, which looks for environment variables (AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY) or a credentials
   * file on the system running this program.
   */
  private static AWSCredentialsProvider getAWSCredentials() {
    return DefaultAWSCredentialsProviderChain.getInstance();
  }
}

KVSUtils类实现

public final class KVSUtils
{
    private static final Logger logger = LoggerFactory.getLogger(KVSUtils.class);

    public enum TrackName {
        AUDIO_FROM_CUSTOMER("AUDIO_FROM_CUSTOMER"),
        AUDIO_TO_CUSTOMER("AUDIO_TO_CUSTOMER");

        private String name;

        TrackName(String name) {
            this.name = name;
        }

        public String getName() {
            return name;
        }
    }

    public static class AudioChunkTrack
    {
        ByteBuffer audio;
        String track;
        long trackNum;

        public AudioChunkTrack(ByteBuffer audio, String track, long trackNum)
        {
            this.audio = audio;
            this.track = track;
            this.trackNum = trackNum;
        }
    }

    /**
     * Fetches the next ByteBuffer of size 1024 bytes from the KVS stream by parsing the frame from the MkvElement
     * Each frame has a ByteBuffer having size 1024
     *
     * @param streamingMkvReader
     * @param fragmentVisitor
     * @param tagProcessor
     * @param contactId
     * @return
     * @throws MkvElementVisitException
     */
    public static AudioChunkTrack getByteBufferFromStream(StreamingMkvReader streamingMkvReader,
                                                     FragmentMetadataVisitor fragmentVisitor,
                                                     KVSContactTagProcessor tagProcessor,
                                                     String contactId) throws MkvElementVisitException {
        while (streamingMkvReader.mightHaveNext()) {
            Optional<MkvElement> mkvElementOptional = streamingMkvReader.nextIfAvailable();
            if (mkvElementOptional.isPresent()) {
                if (tagProcessor.shouldStopProcessing()) {
                    return new AudioChunkTrack(ByteBuffer.allocate(0), null, -1);
                }
                MkvElement mkvElement = mkvElementOptional.get();
                mkvElement.accept(fragmentVisitor);
                if (MkvTypeInfos.SIMPLEBLOCK.equals(mkvElement.getElementMetaData().getTypeInfo())) {
                    MkvDataElement dataElement = (MkvDataElement) mkvElement;
                    @SuppressWarnings("unchecked")
                    Frame frame = ((MkvValue<Frame>) dataElement.getValueCopy()).getVal();
                    ByteBuffer audioBuffer = frame.getFrameData();

                    long trackNumber = frame.getTrackNumber();
                    MkvTrackMetadata metadata = fragmentVisitor.getMkvTrackMetadata(trackNumber);
                    return new AudioChunkTrack(audioBuffer, metadata.getTrackName(), trackNumber);
                }
            }
        }

        return new AudioChunkTrack(ByteBuffer.allocate(0), null, -1);
    }

    private static AudioChunkTrack aggregateChunks(List<AudioChunkTrack> subChunkList, int length, String track)
    {
        // No aggregation if 0 or 1 chunk only.
        if (subChunkList.size() <= 0) return new AudioChunkTrack(ByteBuffer.allocate(0), null, -1);
        if (subChunkList.size() == 1) return subChunkList.get(0);

        // Assume trackNum of 1st chunk is same for all chunks with same track name.
        long trackNum = subChunkList.get(0).trackNum;

        ByteBuffer combinedByteBuffer = ByteBuffer.allocate(length);
        for (AudioChunkTrack sub : subChunkList) {
            combinedByteBuffer.put(sub.audio);
        }
        AudioChunkTrack agg = new AudioChunkTrack(combinedByteBuffer, track, trackNum);
        agg.audio.flip();
        return agg;
    }

    /**
     * Fetches ByteBuffer of provided size from the KVS stream by repeatedly calling {@link KVSUtils#getByteBufferFromStream}
     * and concatenating the ByteBuffers to create a single chunk
     *
     * @param streamingMkvReader
     * @param fragmentVisitor
     * @param tagProcessor
     * @param contactId
     * @param chunkSizeInKB
     * @return
     * @throws MkvElementVisitException
     */
    public static List<AudioChunkTrack> getByteBufferFromStream(StreamingMkvReader streamingMkvReader,
                                                     FragmentMetadataVisitor fragmentVisitor,
                                                     KVSContactTagProcessor tagProcessor,
                                                     String contactId,
                                                     int chunkSizeInKB) throws MkvElementVisitException {

        List<AudioChunkTrack> chunkList = new ArrayList<>();

        for (int i = 0; i < chunkSizeInKB; i++) {
            AudioChunkTrack chunk = KVSUtils.getByteBufferFromStream(streamingMkvReader, fragmentVisitor, tagProcessor, contactId);
            if (chunk.audio.remaining() > 0) {
                chunkList.add(chunk);
            } else {
                break;
            }
        }
        if (chunkList.size() <= 0) return chunkList;

        // Collate sequential chunks by track.
        String lastTrack = null;
        int length = 0, totalLength = 0;
        List<AudioChunkTrack> subChunkList = new ArrayList<>();
        List<AudioChunkTrack> collatedChunks = new ArrayList<>();

        for (AudioChunkTrack act : chunkList) {
            if (lastTrack != null && !lastTrack.equals(act.track)) {
                // Aggregate if there is any change of track.
                collatedChunks.add(aggregateChunks(subChunkList, length, lastTrack));
                subChunkList.clear();
                length = 0;
            }
            subChunkList.add(act);
            length += act.audio.remaining();
            lastTrack = act.track;
            totalLength += length;
        }
        if (subChunkList.size() > 0 && length > 0) {
            // Aggregate any leftover, ungrouped chunks.
            collatedChunks.add(aggregateChunks(subChunkList, length, lastTrack));
        }

        if (totalLength <= 0) {
            return new ArrayList<AudioChunkTrack>();
        }
        return collatedChunks;
    }

    /**
     * Makes a GetMedia call to KVS and retrieves the InputStream corresponding to the given streamName and startFragmentNum
     *
     * @param streamName
     * @param region
     * @param startFragmentNum
     * @param awsCredentialsProvider
     * @return
     */
    public static InputStream getInputStreamFromKVS(String streamName,
                                                    Regions region,
                                                    String startFragmentNum,
                                                    AWSCredentialsProvider awsCredentialsProvider,
                                                    String startSelectorType) {
        Validate.notNull(streamName);
        Validate.notNull(region);
        Validate.notNull(startFragmentNum);
        Validate.notNull(awsCredentialsProvider);

        AmazonKinesisVideo amazonKinesisVideo = (AmazonKinesisVideo) AmazonKinesisVideoClientBuilder.standard().build();

        String endPoint = amazonKinesisVideo.getDataEndpoint(new GetDataEndpointRequest()
                .withAPIName(APIName.GET_MEDIA)
                .withStreamName(streamName)).getDataEndpoint();

        AmazonKinesisVideoMediaClientBuilder amazonKinesisVideoMediaClientBuilder = AmazonKinesisVideoMediaClientBuilder.standard()
                .withEndpointConfiguration(new AwsClientBuilder.EndpointConfiguration(endPoint, region.getName()))
                .withCredentials(awsCredentialsProvider);
        AmazonKinesisVideoMedia amazonKinesisVideoMedia = amazonKinesisVideoMediaClientBuilder.build();

        StartSelector startSelector;
        startSelectorType = isNullOrEmpty(startSelectorType) ? "NOW" : startSelectorType;
        switch (startSelectorType) {
            case "FRAGMENT_NUMBER":
                startSelector = new StartSelector()
                        .withStartSelectorType(StartSelectorType.FRAGMENT_NUMBER)
                        .withAfterFragmentNumber(startFragmentNum);
                logger.info("StartSelector set to FRAGMENT_NUMBER: " + startFragmentNum);
                break;
            case "NOW":
            default:
                startSelector = new StartSelector()
                        .withStartSelectorType(StartSelectorType.NOW);
                logger.info("StartSelector set to NOW");
                break;
        }

        GetMediaResult getMediaResult = amazonKinesisVideoMedia.getMedia(new GetMediaRequest()
                .withStreamName(streamName)
                .withStartSelector(startSelector));

        logger.info("GetMedia called on stream {} response {} requestId {}", streamName,
                getMediaResult.getSdkHttpMetadata().getHttpStatusCode(),
                getMediaResult.getSdkResponseMetadata().getRequestId());

        return getMediaResult.getPayload();
    }
}

核心疑问

我的逻辑未按轨道分离处理,仅读取KVS下一个可用音频块转发至WebSocket,由接收端自行分离轨道。但最终生成的音频嘈杂卡顿、内容缺失,请问问题出在哪里?


内容的提问来源于stack exchange,提问作者Kode Charlie

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 01:00:53